Async processing qua message queue: vì sao đẩy việc nặng ra khỏi request path, và cái giá phải trả bằng eventual consistency
Async processing là mô hình tách một request thành hai giai đoạn: request handler nhận việc, xác nhận với client, rồi giao phần xử lý thật cho một worker chạy ngoài request path — thường qua một message queue (RabbitMQ, AWS SQS, Kafka, Redis Streams, hoặc queue trên nền Redis như BullMQ/Sidekiq). Lý do dev gặp nó trong việc thật rất cụ thể: một endpoint gọi payment provider mất 3s, gửi email confirm mất 1s, resize ảnh mất 5s — nếu làm tuần tự trong request, p99 latency của endpoint là tổng các con số đó, và một downstream chậm hoặc chết đủ để làm timeout hết thread pool của app server. Đẩy vào queue thì request trả về trong vài chục ms; nhưng đổi lại, cái "xong" mà client thấy không còn nghĩa là việc đã thực sự hoàn thành.
Cơ chế hoạt động
Ba thành phần: producer (thường là API server) đóng gói việc thành message rồi publish vào broker; broker (RabbitMQ/SQS/Kafka…) giữ message trong queue có persistence tuỳ cấu hình; consumer/worker poll hoặc được push message, xử lý, rồi ack để broker biết xoá. Nếu worker chết trước khi ack, broker redeliver — đây là gốc của semantic at-least-once: mỗi message được giao ít nhất một lần, có thể nhiều lần. Exactly-once trong hệ phân tán chỉ đạt được ở lớp application bằng cách consumer viết idempotent, không phải bằng cấu hình broker.
Ví dụ với RabbitMQ + Node (amqplib):
// producer — trong HTTP handler
const ch = await conn.createConfirmChannel()
await ch.assertQueue('image.resize', { durable: true })
app.post('/upload', async (req, res) => {
const jobId = crypto.randomUUID()
const payload = Buffer.from(JSON.stringify({ jobId, s3Key: req.body.key }))
await ch.sendToQueue('image.resize', payload, {
persistent: true, // ghi xuống disk, sống sót broker restart
messageId: jobId, // để consumer dedupe
contentType: 'application/json',
})
// đợi broker confirm đã nhận (publisher confirm)
await ch.waitForConfirms()
res.status(202).json({ jobId }) // 202 Accepted — chưa xong, đã nhận
})
// consumer — process riêng
const ch = await conn.createChannel()
await ch.prefetch(8) // chỉ giữ 8 message unacked mỗi consumer
ch.consume('image.resize', async (msg) => {
try {
const job = JSON.parse(msg.content.toString())
await resizeAndUpload(job) // việc thật
ch.ack(msg)
} catch (err) {
// requeue = false → gửi sang dead-letter exchange đã config
ch.nack(msg, false, false)
}
})
Ba chi tiết dễ bỏ sót nhưng quyết định correctness trong production: persistent: true + queue durable: true để message sống sót broker restart; publisher confirm để producer biết broker đã nhận (không có nó thì sendToQueue là fire-and-forget); prefetch để giới hạn số message một consumer giữ, nếu không một consumer chậm sẽ ôm cả queue.
Ngoài "job queue" cổ điển (một producer → một queue → worker pool cạnh tranh), còn hai mô hình hay dùng: pub/sub fanout (một message được nhân bản tới nhiều queue, mỗi consumer group xử lý độc lập — RabbitMQ fanout exchange, Kafka topic với nhiều consumer group) để cùng một event trigger nhiều side-effect độc lập, và log-based queue (Kafka, Redis Streams) nơi message không bị xoá sau khi ack mà tồn theo offset, cho phép replay.
Vấn đề gặp trong production
Failure mode 1: xử lý đồng bộ trong request path. Endpoint đăng ký tài khoản gọi tuần tự: ghi DB, gọi Stripe tạo customer, gửi welcome email qua SendGrid, đẩy analytics event. Tổng latency là tổng của bốn call, và bất kỳ call nào chậm/timeout đều làm cả request fail dù DB đã ghi xong:
// sai — mọi side-effect chặn response
app.post('/signup', async (req, res) => {
const user = await db.insertUser(req.body)
await stripe.customers.create({ email: user.email }) // 500ms–timeout
await sendgrid.send(welcomeEmail(user)) // vài trăm ms
await analytics.track('signup', user) // vài chục ms
res.json(user) // p99 = tổng, và fail-any-fail-all
})
Hậu quả: khi SendGrid chậm, endpoint timeout dù user đã được tạo — client retry, sinh duplicate user; app server thread pool bị các request treo giữ, các endpoint khác cũng chậm theo. Bản đúng là cắt bốn công việc thành bốn message, request chỉ ghi DB + enqueue rồi trả về: