Thực hành: đưa tác vụ AI chạy lâu vào hàng đợi RabbitMQ và trả HTTP 202 ngay
Một request phân tích tài liệu kéo dài vài phút không được giữ kết nối HTTP của khách hàng, và cũng không được biến mất khi worker sập giữa chừng.
- 1Client gửi POST /jobsKèm Idempotency-Key để retry không sinh việc trùng
- 2API trả HTTP 202Header Location trỏ tới URL trạng thái, Retry-After gợi ý thời gian chờ
- 3Message vào queue durableQueue durable, message persistent nên sống qua crash hoặc restart
- 4Worker xử lý, prefetch 1Chỉ ack khi xong; worker chết thì message quay lại queue
- 5Trạng thái hoặc DLQThành công thì cập nhật status; lỗi thì ghi Failed và chuyển sang dead-letter queue
- 6Client poll hoặc nhận webhookĐọc status, percentComplete, error; lần theo correlation ID trong log
Client nhận 202 ngay, còn tác vụ chỉ rời khỏi queue khi worker đã xử lý xong và ack.
Đồ hoạ: FDE Times
Tóm tắt nhanh
- Queue và message phải khai báo durable/persistent, nếu không RabbitMQ quên sạch khi crash hoặc restart.
- Worker chỉ ack khi đã xử lý xong, và đặt prefetch_count=1 để tác vụ dài không dồn hết vào một máy.
- API nhận việc trả 202 ngay, kèm Location để poll và Idempotency-Key để retry không sinh việc trùng.
RabbitMQ có một mặc định dễ gây bất ngờ: khi broker thoát hoặc crash, nó quên luôn các queue và message, trừ khi bạn dặn nó đừng quên. Với một tác vụ AI chạy năm phút, mặc định đó có nghĩa là khách hàng bấm “Phân tích”, chờ, và không bao giờ nhận được gì.
Phần lớn demo LLM gọi model ngay trong request HTTP rồi giữ kết nối cho tới khi có kết quả. Cách đó chạy được trên laptop, nhưng sang môi trường của khách hàng thì gặp timeout của load balancer, client retry rồi sinh việc trùng, còn worker restart thì làm mất việc.
Bài này dựng một pipeline nhỏ theo thứ tự hợp lý khi làm tại hiện trường: queue bền trước, worker an toàn sau, API bất đồng bộ cuối cùng. Khoảng một buổi tối là xong, và bạn có thứ để mang vào phỏng vấn.
Bạn sẽ dựng cái gì?
AWS định nghĩa hàng đợi thông điệp là một hình thức giao tiếp bất đồng bộ giữa các dịch vụ, dùng để tách các tác vụ xử lý nặng khỏi phần còn lại của ứng dụng. Gọi model chính là loại tác vụ nặng đó.
Luồng cần dựng gồm bốn mảnh. Client gửi POST /jobs và nhận ngay HTTP 202. API đóng gói tác vụ thành message, đẩy vào queue ai_jobs. Worker lấy message ra, gọi model, cập nhật trạng thái, còn client poll GET /jobs/{id} để xem kết quả.
Bạn cần Python 3, thư viện pika (client mà tutorial chính thức của RabbitMQ dùng) và một RabbitMQ chạy ở localhost, cài theo hướng dẫn chính thức. Các đoạn code dưới đây đã được rút gọn để dạy ý tưởng: chưa có cấu hình kết nối, chưa xử lý reconnect, còn trạng thái thì lưu trong bộ nhớ.
Bước 1: Một queue sống sót qua lần restart
Tutorial Work Queues của RabbitMQ mô tả đúng bài toán này: thay vì chạy ngay một tác vụ tốn tài nguyên rồi ngồi chờ, ta đóng gói nó thành message và gửi vào queue. Muốn message sống qua crash, cả queue lẫn message đều phải được đánh dấu bền.
import json, pika
conn = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
ch = conn.channel()
ch.queue_declare(queue='ai_jobs', durable=True)
def enqueue(job: dict):
ch.basic_publish(
exchange='',
routing_key='ai_jobs',
body=json.dumps(job),
properties=pika.BasicProperties(
delivery_mode=pika.DeliveryMode.Persistent),
)
Kiểm tra: gửi vài job, restart RabbitMQ rồi xem queue còn message hay không. Nếu queue trống, bạn đã quên durable=True hoặc quên delivery_mode ở một trong hai chỗ.
Bước 2: Worker chỉ ack khi đã xong việc
Đây là chỗ bảo vệ tác vụ dài. Nếu worker chết trước khi ack, RabbitMQ hiểu rằng message chưa được xử lý trọn vẹn và đưa nó trở lại queue. Muốn vậy bạn phải ack thủ công, và gọi ack ở dòng cuối cùng chứ không phải dòng đầu.
ch.basic_qos(prefetch_count=1)
def handle(ch, method, props, body):
job = json.loads(body)
set_status(job['id'], 'Running')
result = run_model(job) # gọi LLM, có thể mất vài phút
set_status(job['id'], 'Succeeded', result=result)
ch.basic_ack(delivery_tag=method.delivery_tag)
ch.basic_consume(queue='ai_jobs', on_message_callback=handle)
ch.start_consuming()
Handler này chưa xử lý ngoại lệ: nếu run_model ném lỗi, job không bao giờ được đánh dấu Failed. Phần về dead-letter queue ở cuối bài sẽ bổ sung đoạn đó.
Dòng basic_qos(prefetch_count=1) bảo RabbitMQ không giao cho một worker quá một message mỗi lần. Thử hình dung hai worker và một dãy job xen kẽ: job tóm tắt hợp đồng mất 10 phút, job phân loại email mất 10 giây.
Nếu chia đều theo lượt mà không có prefetch, toàn bộ job 10 phút có thể rơi vào cùng một worker trong khi worker kia ngồi không. Có prefetch bằng 1, ai rảnh trước thì nhận việc trước, và hàng đợi tự cân tải theo độ dài thật của từng tác vụ.
Kiểm tra: cho run_model ngủ 60 giây, kill worker ở giây thứ 30 rồi khởi động lại. Job phải được chạy lại từ đầu. Hệ quả kèm theo là run_model cần an toàn khi chạy lại, vì một job có thể được xử lý hai lần.
Bước 3: API trả 202, không bắt client phải chờ
Mẫu Async Request-Reply trong Azure Architecture Center mô tả phần nhìn ra ngoài: API trả HTTP 202 (Accepted) ngay để xác nhận đã nhận yêu cầu. Phản hồi kèm header Location trỏ tới URL mà client poll trạng thái, và Retry-After gợi ý nên chờ bao lâu trước lần poll kế tiếp.
Client nào cũng sẽ retry khi mạng chập chờn. Vì thế Microsoft khuyên dùng Idempotency-Key: nếu backend nhận một key đã thấy, nó trả lại resource trạng thái sẵn có chứ không đẩy thêm một work item thứ hai vào queue. Phác thảo dưới đây không gắn với framework nào, và từ chối luôn request thiếu key:
JOBS, KEYS = {}, {} # rút gọn: production dùng database
def post_jobs(request):
key = request.headers.get('Idempotency-Key')
if not key:
return 400, {'error': 'Idempotency-Key is required'}
if key in KEYS:
return accepted(KEYS[key])
job_id = new_id()
JOBS[job_id] = {'status': 'Pending', 'createdAt': now(),
'lastUpdatedAt': now()}
KEYS[key] = job_id
enqueue({'id': job_id, 'input': request.json,
'correlationId': job_id})
return accepted(job_id)
def accepted(job_id):
return 202, {'Location': f'/jobs/{job_id}', 'Retry-After': '10'}
Kiểm tra: gửi hai POST cùng key, queue chỉ được tăng đúng một message và cả hai phản hồi phải trỏ về cùng một Location. Gửi thêm một POST không có header: bạn phải nhận 400 và queue không đổi.
Bước 4: Endpoint trạng thái cho client biết chính xác job đang ở đâu
GET /jobs/{id} trả về những trường mà mẫu Async Request-Reply gợi ý: status (Pending, Running, Succeeded, Failed hoặc Canceled), createdAt, lastUpdatedAt, percentComplete và error. Trong bản demo, việc này chỉ là đọc JOBS[job_id].
Poll không phải lựa chọn duy nhất. Nordic APIs và Postman đều mô tả trường hợp server báo cho client qua callback, chẳng hạn webhook. Một cách kết hợp hợp lý là dùng webhook cho hệ thống của khách và giữ endpoint trạng thái cho người cần mở ra kiểm tra.
Lỗi xảy ra thì message đi đâu?
Azure liệt kê mất dữ liệu là một thách thức của kiến trúc bất đồng bộ, và cách xử lý là lưu bền sự kiện đang trên đường đi, chỉ dequeue khi thành phần kế tiếp đã ack. Bước 1 và 2 đã làm đúng điều đó. Còn lại câu hỏi: job lỗi thật thì sao?
Gợi ý của Azure là chuyển sự kiện lỗi sang dead-letter queue (DLQ) để quản trị viên kiểm tra. Bản rút gọn dưới đây thay handler ở Bước 2 và chỉ dùng những API đã có: bọc run_model trong try, khi lỗi thì ghi Failed kèm error, publish message vào queue ai_jobs_dlq khai báo durable, rồi mới ack message gốc.
ch.queue_declare(queue='ai_jobs_dlq', durable=True)
def handle(ch, method, props, body):
job = json.loads(body)
set_status(job['id'], 'Running')
try:
result = run_model(job)
set_status(job['id'], 'Succeeded', result=result)
except Exception as e: # rút gọn: bắt mọi loại lỗi
set_status(job['id'], 'Failed', error=str(e))
ch.basic_publish(
exchange='',
routing_key='ai_jobs_dlq',
body=body,
properties=pika.BasicProperties(
delivery_mode=pika.DeliveryMode.Persistent),
)
ch.basic_ack(delivery_tag=method.delivery_tag)
Đây là bản giản lược: mọi lỗi, kể cả lỗi mạng thoáng qua, đều đi thẳng vào DLQ. Kiểm tra: cho run_model ném lỗi với một input cụ thể, rồi xác nhận GET /jobs/{id} trả Failed, ai_jobs_dlq có thêm một message và ai_jobs không còn giữ job đó.
Mảnh cuối là correlation ID. Azure khuyên gắn nó vào mọi sự kiện để mọi consumer và hệ thống log nối được các thao tác liên quan thành một trace. Bản demo dùng luôn job_id, nên bạn hãy in nó ở mọi dòng log của API và worker, rồi thử grep một ID để thấy toàn bộ hành trình của job đó.
Bốn lỗi hay gặp nhất
Lỗi phổ biến nhất là ack ngay khi vừa nhận message “cho gọn”. Worker chết giữa chừng thì job mất mà không để lại dấu vết, vì RabbitMQ tin rằng việc đã xong.
Lỗi thứ hai là khai báo queue durable nhưng quên đánh dấu message persistent, hoặc làm ngược lại. Phải có đủ cả hai thì mới qua được một lần restart.
Lỗi thứ ba là bỏ qua Idempotency-Key vì nghĩ client sẽ không retry. Với tác vụ AI, mỗi job trùng là thêm một lần gọi model tốn tiền, và có thể sinh ra kết quả thứ hai mâu thuẫn với kết quả đầu.
Lỗi thứ tư là khai báo trường percentComplete nhưng worker không bao giờ cập nhật nó. Client thấy một job 10 phút đứng ở 0% suốt 9 phút và kết luận hệ thống đã treo. Nếu tác vụ chia được thành các bước, chẳng hạn từng trang tài liệu, hãy cập nhật percentComplete và lastUpdatedAt sau mỗi bước.
Tại hiện trường và trong CV
Nếu khách hàng phàn nàn rằng “AI hay bị treo”, một cách tiếp cận hợp lý là dựng chính bộ khung trên trước khi tinh chỉnh model. Nên hỏi trước tiên: job dài nhất mất bao lâu, ai đang retry, và nếu worker restart lúc 2 giờ sáng thì job đang chạy sẽ đi đâu.
Khi đọc JD, bạn có thể để ý các cụm như “asynchronous processing”, “message queue”, “idempotency” hay “observability”. Đó là tín hiệu nên hỏi kỹ hơn trong phỏng vấn, chứ chưa phải bằng chứng chắc chắn về công việc. Trên CV, đừng chỉ ghi “dùng RabbitMQ”.
Hãy viết rằng bạn đã thiết kế API 202 có Idempotency-Key, worker có manual ack và prefetch, có DLQ và correlation ID, và đã kiểm chứng bằng cách kill worker giữa job.
Model sẽ còn thay đổi nhiều lần. Những thứ đã dựng ở trên, gồm queue bền, ack đúng lúc và một API luôn cho client biết job đang ở bước nào, vẫn dùng được ở mọi dự án tiếp theo.
6 nguồn
- What is a Message Queue? (AWS)
- RabbitMQ tutorial - Work Queues | RabbitMQ
- Asynchronous Request-Reply Pattern - Azure Architecture Center | Microsoft Learn · 2026-03-30
- Event-Driven Architecture Style - Azure Architecture Center | Microsoft Learn · 2026-03-06
- The Differences Between Synchronous and Asynchronous APIs · 2024-01-25
- Understanding asynchronous APIs · 2022-08-26