# Thực hành: dựng pipeline xử lý hồ sơ qua hàng đợi, kết hợp pipes and filters với claim check

> Một file PDF bị nhét thẳng vào hàng đợi, một bước xử lý chạy lại hai lần, một client gửi lại POST: ba lỗi quen thuộc này đều xử lý được bằng vài chục dòng Python, nếu bạn thiết kế đúng từ đầu.

Bản gốc: https://fdetimes.net/vi/bach-khoa/thuc-hanh-pipes-filters-claim-check-xu-ly-tai-lieu/

Chi tiết dễ bị bỏ qua nhất trong tài liệu mẫu của Microsoft về pipes and filters là message chạy giữa các bước không chứa bức ảnh cần xử lý. Message chỉ chứa một tham chiếu (claim check) trỏ tới ảnh đang nằm trong storage.

Cách làm này giữ cho message nhỏ gọn, nhưng không giải quyết được một rủi ro khác: hàng đợi có thể giao cùng một message hơn một lần. Vì thế, ngoài claim check, mỗi filter còn phải idempotent, tức là nhận một message hai lần vẫn cho ra đúng một kết quả. Bước 3 sẽ xử lý chuyện này.

Ở dự án của FDE, bạn sẽ gặp đúng bài toán này dưới nhiều tên gọi: hồ sơ vay, hợp đồng scan, đơn bảo hiểm. Khách hàng gửi lên một file nặng, file phải đi qua vài bước kiểm tra và trích xuất, rồi ai đó phải biết khi nào việc xử lý xong. Bài này dựng một pipeline như vậy trên laptop, chỉ dùng thư viện chuẩn của Python.

## Bạn sẽ dựng gì, và cần những gì?

Pipeline gồm ba bước validate, extract và store, nối với nhau bằng hàng đợi. Đây chính là cách tài liệu best practices về background job của Microsoft gợi ý: một job đi qua nhiều giai đoạn thì tách thành các filter rời nhau, nối bằng queue. Bên ngoài pipeline là một API nhận hồ sơ và trả trạng thái xử lý.

Bạn chỉ cần Python 3 và một thư mục trống. Toàn bộ code dưới đây là **mô phỏng giản lược**: `queue.Queue` đóng vai message broker, một thư mục local đóng vai blob storage, và API viết dưới dạng hàm thuần chứ không dùng framework nào. Khi làm ở chỗ khách hàng, bạn thay từng phần bằng dịch vụ thật của họ, còn logic giữ nguyên.

## Bước 1: ghi file trước, gửi phiếu sau

Theo mô tả của Microsoft, claim check là cách cất payload lớn vào một kho bên ngoài và chỉ gửi tham chiếu qua hệ thống messaging. Consumer nhận tham chiếu rồi tự lấy payload về. Thứ tự hai thao tác này rất quan trọng: tài liệu yêu cầu chỉ publish tham chiếu khi việc ghi payload đã thành công.

```python
import hashlib, queue, uuid
from pathlib import Path

STORE = Path("blobs"); STORE.mkdir(exist_ok=True)
q_validate, q_extract, q_store, q_dead = (queue.Queue() for _ in range(4))

def put_payload(data: bytes) -> str:
key = hashlib.sha256(data).hexdigest()
(STORE / key).write_bytes(data)
return key

def submit(data: bytes) -> str:
job_id = str(uuid.uuid4())
ref = put_payload(data)                      # 1. ghi payload
q_validate.put({"job_id": job_id, "blob": ref, "schema": "hoso.v1"})  # 2. publish
return job_id
```

**Kiểm tra:** gọi `submit(b"%PDF-1.4 ...")` rồi xem thư mục `blobs/`. Trong đó phải có đúng một file, còn message trong `q_validate` chỉ gồm vài chục byte.

Dùng hash của nội dung làm key là một lựa chọn có chủ đích: cùng một file thì luôn ra cùng một key. Nhờ vậy, nếu bước ghi bị retry thì cũng không sinh thêm bản sao.

Microsoft lưu ý rằng ghi và publish không xảy ra nguyên tử, nên bạn phải tính tới cả message trùng lẫn payload mồ côi, tức file đã ghi nhưng message không bao giờ được gửi đi.

## Bước 2: viết khuôn chung cho filter, để pipe chỉ chuyển tin

Tài liệu Azure nói rõ: pipe không định tuyến và không chứa logic, nó chỉ đưa output của filter này thành input của filter kế tiếp. Vì thế mọi xử lý, kể cả xử lý lỗi, đều nằm trong filter.

```python
def run_filter(inbox, outbox, work, max_attempts=3):
msg = inbox.get()
try:
result = work(msg)
if outbox is not None:
outbox.put({**msg, **result})
except Exception:
msg["attempts"] = msg.get("attempts", 0) + 1
(q_dead if msg["attempts"] >= max_attempts else inbox).put(msg)
```

Nhánh `except` là chỗ xử lý poison message. Một message lỗi lặp đi lặp lại sẽ không bị đưa vào hàng đợi mãi mãi mà được chuyển sang dead-letter queue, đúng như khuyến nghị trong tài liệu background job.

Trường `schema` gắn trong message cũng có lý do: các filter chỉ biết schema vào và ra của mình, nên khi schema được chuẩn hóa, bạn có thể đổi thứ tự các bước mà không phải viết lại chúng.

## Bước 3: ba filter, và filter cuối phải idempotent

```python
DONE = {}

def validate(msg):
data = (STORE / msg["blob"]).read_bytes()
if not data.startswith(b"%PDF"):
raise ValueError("không phải PDF")
return {}

def extract(msg):
data = (STORE / msg["blob"]).read_bytes()
return {"size_bytes": len(data)}   # giản lược: thay bằng OCR/trích xuất thật

def store(msg):
if msg["job_id"] in DONE:          # đã xử lý thì bỏ qua
return {}
DONE[msg["job_id"]] = {"blob": msg["blob"], "size_bytes": msg["size_bytes"]}
return {}
```

Lý do của dòng `if` trong `store`: hàng đợi chỉ đảm bảo giao tin *ít nhất một lần*, nghĩa là cùng một message có thể đến nhiều lần.

Microsoft mô tả thêm một tình huống cụ thể: filter đã post kết quả rồi mới hỏng, message được chạy lại trên một instance khác, và kết quả bị nhân đôi.

Vì vậy filter nào ghi ra thế giới bên ngoài cũng cần một khóa để nhận biết việc đã làm.

**Kiểm tra:** chạy lần lượt ba filter bằng `run_filter(q_validate, q_extract, validate)`, `run_filter(q_extract, q_store, extract)` và `run_filter(q_store, None, store)`. Sau đó đưa thủ công đúng message đó vào `q_store` thêm lần nữa rồi chạy lại. `len(DONE)` phải vẫn bằng 1.

**Điểm mấu chốt:** Ghi dữ liệu trước rồi mới gửi phiếu, và filter nào cũng phải chịu được việc nhận một tin hai lần.

## Bước 4: tìm filter chậm nhất trước khi tối ưu

Tài liệu Azure viết rằng thời gian xử lý một request phụ thuộc vào những filter chậm nhất trong pipeline. Thử hình dung validate mất 1 giây, extract mất 8 giây và store mất 1 giây.

Tối ưu validate gần như không thay đổi gì. Cách đúng là chạy song song nhiều instance của extract cùng đọc từ `q_extract`, và đây là lợi ích trực tiếp của việc tách các bước ra bằng hàng đợi.

Trong bản mô phỏng, bạn có thể tạo vài `threading.Thread` cùng gọi `run_filter` trên `q_extract`. Vì `store` đã idempotent nên dù hai worker có lỡ xử lý trùng một message, kết quả vẫn không bị sai.

## Bước 5: API trả 202, chặn gửi trùng, xong thì trả 303

Khách hàng không nên phải giữ kết nối trong lúc pipeline chạy. Theo mẫu async request-reply, API kiểm tra request rồi trả ngay HTTP 202 kèm header `Location` và `Retry-After`, sau đó client tự poll endpoint trạng thái.

```python
JOBS_BY_KEY = {}

def post_hoso(body: bytes, idem_key: str):
job_id = JOBS_BY_KEY.get(idem_key)
if job_id is None:
job_id = submit(body)
JOBS_BY_KEY[idem_key] = job_id
return 202, {"Location": f"/hoso/status/{job_id}", "Retry-After": "5"}

def get_status(job_id: str):
if job_id in DONE:
return 303, {"Location": f"/hoso/{job_id}"}
return 200, {"status": "processing"}
```

Hai chi tiết trong đoạn code này chặn hai sự cố khác nhau. Header `Idempotency-Key` cho phép backend trả về status resource đã có khi client gửi trùng, thay vì đưa thêm một work item nữa vào hàng đợi.

Còn khi job xong, Microsoft khuyên dùng 303 See Other thay cho 302, vì với một số client, 302 có thể khiến POST ban đầu bị gửi lại.

**Kiểm tra:** gọi `post_hoso` hai lần với cùng một key. Cả hai lần phải trả về cùng một `Location`, và `q_validate` chỉ được tăng thêm một message.

## Bài tập cuối: cố tình làm validate hỏng

Bài tập này ráp cả năm bước lại với nhau. Gọi `post_hoso(b"hello", "key-loi")`: file vẫn được ghi vào `blobs/`, API vẫn trả 202, nhưng nội dung không bắt đầu bằng `%PDF` nên `validate` sẽ ném lỗi.

Giờ chạy `run_filter(q_validate, q_extract, validate)` ba lần. Hai lần đầu, message quay lại `q_validate` với `attempts` tăng dần; lần thứ ba nó rơi vào `q_dead`. Kiểm tra `q_dead.qsize()` bằng 1, `q_extract` vẫn rỗng, và `get_status` của job này vẫn trả `processing`.

Kết quả cuối cùng cho thấy một lỗ hổng của bản mô phỏng: client sẽ poll mãi mà không biết hồ sơ đã hỏng. Bước tiếp theo nên làm là cho endpoint trạng thái đọc cả dead-letter queue và trả về trạng thái lỗi.

## Những lỗi thường gặp

Lỗi đầu tiên là publish message trước khi ghi xong file, ngược với thứ tự mà tài liệu claim check yêu cầu. Khi đó filter đầu tiên có thể đọc phải một tham chiếu trỏ vào khoảng trống.

Lỗi thứ hai là không ai chịu trách nhiệm xóa payload: tài liệu claim check yêu cầu chỉ định rõ ai sở hữu việc xóa và lưu giữ, đồng thời khớp thời hạn lưu payload với thời hạn sống của message, để một tham chiếu còn hợp lệ không trỏ vào dữ liệu đã bị xóa.

Lỗi thứ ba liên quan tới bảo mật: nhét security token vào claim check cho tiện. Microsoft khuyên không làm như vậy. Lỗi cuối cùng là áp dụng pattern này sai chỗ. Nếu việc xử lý bắt buộc phải xong trong chính request ban đầu, kiểu request/response, thì pipes and filters không phù hợp.

## Kỹ năng này xuất hiện thế nào ở chỗ khách hàng?

Ở site khách hàng, việc đầu tiên không phải là viết code mà là hỏi một câu: người dùng có chấp nhận nhận kết quả sau, qua một endpoint trạng thái, hay họ bắt buộc phải có kết quả ngay?

Câu trả lời quyết định bạn có dùng pipeline hay không. Tiếp theo, hãy vẽ các bước hiện tại thành filter, đánh dấu bước chậm nhất và hỏi xem file đang được lưu ở đâu, ai được phép xóa.

Trong CV, đừng chỉ ghi "dùng message queue". Hãy viết rõ bạn đã xử lý giao tin trùng thế nào, đặt dead-letter queue ở đâu và ai chịu trách nhiệm dọn payload.

Cuốn Enterprise Integration Patterns đặt ra cho pattern này câu hỏi: làm sao xử lý một message qua nhiều bước phức tạp mà các bước vẫn độc lập và linh hoạt. Bài thực hành vừa rồi cho thấy câu trả lời: mỗi filter chỉ biết schema vào và ra của riêng nó, nên bạn đổi thứ tự hay chạy thêm instance mà không phải viết lại filter nào.

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

- Chạy đoạn mã mô phỏng trong bài, đẩy cùng một message vào hàng đợi store hai lần rồi kiểm tra xem kết quả có bị nhân đôi không.
- Mở một hệ thống batch bạn đang bảo trì, vẽ lại thành các filter và ghi rõ bước nào chậm nhất. Đó là chỗ đầu tiên nên chạy thêm instance.
- Thêm vào CV một dòng mô tả pipeline bạn từng làm, trong đó nêu rõ cách bạn xử lý giao tin trùng và dead-letter queue.

## Nguồn

- [Pipes and Filters pattern - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/patterns/pipes-and-filters)

- [Claim Check pattern - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/patterns/claim-check)

- [Best Practices for Background Jobs - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/best-practices/background-jobs)

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

- [Pipes and Filters - Enterprise Integration Patterns](http://www.enterpriseintegrationpatterns.com/patterns/messaging/PipesAndFilters.html)
