FDE PulseViệc làm FDE đang mở 316Mới đăng 7 ngày qua 10Chủ đề nổi bật: Đào tạo kỹ năng FDE tại Đông Nam Á

Tờ báo của nghề Forward Deployed Engineer

Bách khoa

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.

Đồ hoạMột hồ sơ đi qua pipeline có claim check
  1. 1API nhận hồ sơKiểm tra Idempotency-Key rồi trả ngay 202 kèm Location và Retry-After
  2. 2Ghi file vào storageChỉ khi ghi xong mới publish tham chiếu (claim check) vào hàng đợi đầu tiên
  3. 3Filter validateĐọc file qua tham chiếu; lỗi lặp lại thì chuyển sang dead-letter queue
  4. 4Filter extractThường là bước chậm nhất, nên chạy song song nhiều instance
  5. 5Filter storeIdempotent theo job_id nên nhận trùng message cũng không ghi đôi
  6. 6Endpoint trạng tháiTrả 303 See Other khi hồ sơ đã xử lý xong

Chỉ có tham chiếu tới file đi qua các hàng đợi, còn bản thân file nằm yên trong storage từ đầu đến cuối.

Đồ hoạ: FDE Times

Tóm tắt nhanh

  • Pipe chỉ chuyển message đi tiếp. Mọi logic nằm trong filter, và mỗi filter chỉ cần biết schema đầu vào và đầu ra của nó.
  • File lớn được ghi vào kho trước, sau đó mới publish tham chiếu. Hai thao tác này không nguyên tử nên filter phải idempotent.
  • Khi khách hàng cần kết quả ngay trong cùng một request, pipeline bất đồng bộ không phải lựa chọn phù hợp.
Chia sẻLinkedInFacebookX

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.

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.

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

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.

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.

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.

5 nguồn
Đọc tiếp trên lộ trình · Chặng 5: Triển khaiBa replica, một cron job: làm sao để khách không nhận ba emailScale out từ một bản chạy lên ba thì cron job cũng chạy ba lần. Scheduler-Agent-Supervisor và leader election giúp kiểm soát chuyện này, nhưng cách sửa bền nhất không nằm ở chiếc khoá.