Khi hệ thống của khách đổ dồn dữ liệu: dùng queue, competing consumers và back-pressure để không sập
Khách sẽ không bao giờ gửi dữ liệu theo tốc độ bạn muốn, nên hệ thống của bạn phải tự quyết định nó nhận nhanh đến đâu.
- 1Webhook nhận và ghi vào queueTrả 200 ngay, không xử lý tại chỗ; queue làm bộ đệm cho đợt dồn
- 2Pub/sub chia bản saoMỗi subscriber (chấm rủi ro, đồng bộ CRM) nhận mọi message
- 3Flow control ở consumerGiới hạn message chưa ack: prefetch, multiplier bằng 1, timeout lớn hơn p99
- 4Competing consumers xử lýMỗi message một worker; consumer idempotent để chịu bản trùng và sai thứ tự
- 5Scale và cô lập lỗiAutoscale theo độ sâu queue; message fail nhiều lần chuyển sang DLQ
Tách nhận khỏi xử lý, rồi để phía consumer tự quyết tốc độ bằng flow control và scale.
Đồ hoạ: FDE Times
Tóm tắt nhanh
- Đặt queue giữa hệ thống khách và worker để hấp thụ đợt tải dồn, rồi để nhiều worker cùng tranh nhau xử lý.
- Back-pressure nằm ở phía consumer: giới hạn message chưa ack, đặt timeout khớp thời gian xử lý thật, scale theo độ sâu queue.
- Thứ tự không được đảm bảo và message có thể đến hai lần, nên consumer phải idempotent và có DLQ.
Thử hình dung: hệ thống ERP của khách chạy job đồng bộ cuối ngày và bắn 50.000 đơn hàng vào webhook của bạn trong mười phút. Mỗi đơn cần gọi một model, ghi database, đẩy kết quả sang CRM. API của bạn bắt đầu timeout, phía khách tự retry, và sáng hôm sau bạn ngồi dọn đơn bị ghi hai lần.
Đây là kiểu tình huống mà ai làm tích hợp với hệ thống của khách cũng nên chuẩn bị trước. Bạn không kiểm soát được hệ thống của khách, không thể bảo họ “gửi chậm lại”. Thứ bạn kiểm soát được là cách hệ thống của mình nhận việc, và kỹ năng đó xoay quanh ba ý tưởng: pub/sub, competing consumers và back-pressure.
Queue không làm bạn nhanh hơn, nó cho bạn quyền chọn tốc độ
Sai lầm ở ví dụ trên là webhook vừa nhận vừa xử lý luôn. Cách sửa là tách làm hai: webhook chỉ ghi message vào một queue rồi trả 200 ngay, còn một nhóm worker kéo message ra xử lý theo nhịp của mình.
Tài liệu kiến trúc của Microsoft mô tả đúng vai trò này: queue là bộ đệm giữa phía gửi và các instance xử lý, san phẳng lưu lượng lúc lên lúc xuống.
Khi nhiều worker cùng đọc một queue, ta có mẫu competing consumers: các worker tranh nhau, mỗi message chỉ một worker nhận. Pub/sub thì ngược lại, mọi subscriber đều nhận mọi message. Khác biệt cốt lõi nằm ở chỗ ai nhận mỗi message, và nó quyết định bạn thiết kế topology thế nào.
Trong thực tế bạn thường cần cả hai. Sự kiện “đơn hàng mới” được publish một lần; subscriber “chấm điểm rủi ro” và subscriber “đồng bộ CRM” đều nhận bản sao riêng. Bên trong mỗi subscriber lại là một nhóm worker cạnh tranh. Lợi ích kèm theo: một subscriber hỏng không kéo publisher hay subscriber khác sập theo, và broker giữ message lại chờ nó hồi phục.
Back-pressure nằm ở đâu?
Queue hấp thụ được đợt dồn, nhưng nếu worker tham lam kéo về quá nhiều message một lúc, áp lực chỉ chuyển từ webhook sang RAM của worker.
Hướng dẫn pub/sub của Microsoft đưa ra thứ tự ưu tiên rõ ràng: trước hết dùng flow control của broker để giới hạn số message chưa ack của mỗi subscriber, rồi mới scale out bằng competing consumers khi flow control không đủ.
Mỗi broker đặt chỗ điều tiết này ở một nơi khác, dưới một tên khác:
| Broker | Chỗ chỉnh back-pressure | Điều cần nhớ |
|---|---|---|
| RabbitMQ | prefetch (basic.qos) |
Giá trị 0 nghĩa là không giới hạn message chưa ack |
| Celery | worker_prefetch_multiplier |
Task chạy lâu thì đặt bằng 1 để worker không ôm việc |
| Amazon SQS | visibility timeout | Mặc định 30 giây; standard queue có trần khoảng 120.000 message in-flight |
| Kafka | mô hình pull + consumer group | Lên kế hoạch số partition từ đầu: mỗi partition chỉ thuộc một consumer trong group, nên số partition là trần song song |
Kafka đáng nói riêng. Vì consumer tự kéo dữ liệu, tài liệu thiết kế của Kafka chỉ ra rằng consumer chậm chỉ đơn giản là tụt lại phía sau rồi bắt kịp khi có thể. Back-pressure gần như có sẵn. Đổi lại, nếu topic có 6 partition thì consumer thứ 7 trong cùng group sẽ ngồi chơi, nên hãy tính mức song song cần thiết khi thiết kế topic.
Một ví dụ làm từ đầu đến cuối
Quay lại 50.000 đơn. Giả sử mỗi đơn mất trung bình 2 giây, chậm nhất 45 giây khi model phản hồi lâu. Tổng khối lượng là 100.000 worker-giây; với 20 worker, queue cạn sau khoảng 5.000 giây, tức hơn 83 phút. Câu hỏi đầu tiên cho khách: kết quả có cần xong trước 83 phút không? Nếu có, bạn biết ngay cần bao nhiêu worker.
Bây giờ đến cái bẫy. Nếu bạn dùng SQS với visibility timeout mặc định 30 giây, đơn nào chạy 45 giây sẽ hiện lại trong queue trước khi worker kịp xoá, và một worker khác nhặt nó lên.
Cơ chế đó được thiết kế để cứu việc của worker bị crash, nhưng ở đây nó tạo ra bản trùng. Hãy đặt timeout cao hơn thời gian xử lý chậm nhất thực đo.
Kể cả khi timeout đúng, AWS vẫn nói rõ SQS giao ít nhất một lần, nên không có bảo đảm tuyệt đối rằng message không đến hai lần. Microsoft cũng lưu ý thứ tự nhận message giữa các consumer không phản ánh thứ tự tạo ra. Vì thế consumer phải idempotent:
def handle(msg):
key = f"{msg['order_id']}:{msg['version']}"
with db.transaction():
if db.exists("processed_keys", key):
return # bản trùng, bỏ qua
if db.current_version(msg["order_id"]) > msg["version"]:
return # bản cũ đến muộn, bỏ qua
upsert_order(msg)
db.insert("processed_keys", key)
Hai điều kiện ở đây giải hai bài toán khác nhau. Khoá idempotency chặn message trùng; so sánh version chặn message cũ đến sau message mới. Nếu hệ thống của khách không gửi version, hãy hỏi họ trường nào tăng dần, ví dụ updated_at, ngay trong buổi discovery.
Bước cuối là message độc: một đơn có dữ liệu hỏng sẽ fail mãi. Đừng để nó quay vòng vô hạn; sau một số lần giao nhất định, đẩy nó sang dead-letter queue để người xem xét sau. Với worker pool, bạn có thể autoscale theo độ sâu queue, kể cả xuống 0 khi queue rỗng, đồng thời đặt trần concurrency để không bóp chết database phía sau.
Tự làm theo thứ tự nào?
Đầu tiên, đo: lưu lượng đỉnh của khách, thời gian xử lý trung bình và chậm nhất, cùng hạn chót nghiệp vụ. Sau đó tách nhận và xử lý bằng một queue, quyết định chỗ nào cần pub/sub (nhiều hệ thống cùng quan tâm một sự kiện) và chỗ nào chỉ cần competing consumers.
Tiếp theo chỉnh các tham số điều tiết: prefetch nhỏ, visibility timeout lớn hơn p99, số partition được tính trước cho mức song song mục tiêu. Cuối cùng viết consumer idempotent, cấu hình DLQ và đặt cảnh báo trên độ sâu queue cùng số message in-flight.
Những lỗi gặp nhiều nhất
Lỗi phổ biến nhất là để prefetch bằng 0 vì tưởng “càng nhiều càng nhanh”, để rồi một worker ôm hàng nghìn message trong khi worker khác rảnh. Kế đến là giữ nguyên visibility timeout mặc định cho tác vụ gọi model chạy lâu. Một lỗi khác là thêm consumer Kafka vượt số partition rồi ngạc nhiên vì throughput không tăng.
Hai lỗi còn lại nguy hiểm hơn vì chúng im lặng: giả định message đến đúng thứ tự, và không có DLQ nên một bản ghi hỏng chặn cả hàng đợi.
Chính chuỗi quyết định trong ví dụ 50.000 đơn là thứ nên mang vào buổi phỏng vấn FDE: lưu lượng đỉnh, phép tính số worker, vì sao timeout phải lớn hơn thời gian chậm nhất, khoá idempotency và chỗ đặt DLQ. Kể lại một lần như vậy bằng số liệu thật của bạn thuyết phục hơn nhiều so với một dòng “có kinh nghiệm event-driven” trong CV.
Lần tới khi khách hỏi “hệ thống của anh chịu được bao nhiêu request mỗi giây?”, câu trả lời tốt nhất không phải một con số, mà là: hệ thống của bạn nhận bao nhiêu cũng được, và xử lý đúng tốc độ nó chịu nổi.