# 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.

Bản gốc: https://fdetimes.net/vi/bach-khoa/thuc-hanh-hang-doi-cho-tac-vu-ai-chay-lau/

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.

```python
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.

```python
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:

```python
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.

```python
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.

**Điểm mấu chốt:** Với tác vụ AI chạy lâu, ack đúng lúc quan trọng hơn chọn đúng model.

## 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.

**Thử ngay tuần này:**

- Dựng bản demo ai_jobs theo bốn bước trong bài, rồi kill worker khi nó đang chạy một tác vụ giả lập 60 giây và xác nhận message quay lại queue.
- Gửi cùng một POST hai lần với cùng Idempotency-Key, rồi gửi thêm một POST không có key: số message trong queue chỉ được tăng 1 và request thiếu key phải bị từ chối.
- Viết một đoạn README cho demo, ghi rõ bạn xử lý mất message, retry trùng và lỗi bằng cách nào, rồi đưa link vào CV.

## Nguồn

- [What is a Message Queue? (AWS)](https://aws.amazon.com/message-queue/)

- [RabbitMQ tutorial - Work Queues | RabbitMQ](https://www.rabbitmq.com/tutorials/tutorial-two-python)

- [Asynchronous Request-Reply Pattern - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/patterns/async-request-reply)

- [Event-Driven Architecture Style - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/guide/architecture-styles/event-driven)

- [The Differences Between Synchronous and Asynchronous APIs](https://nordicapis.com/the-differences-between-synchronous-and-asynchronous-apis/)

- [Understanding asynchronous APIs](https://blog.postman.com/understanding-asynchronous-apis/)
