Skip to content

Pub/Sub + Dataflow: Cơ Chế Xử Lý Chính Xác Một Lần (Exactly-Once)

Why this matters in production

Trong xử lý luồng dữ liệu thời gian thực (stream processing), việc đảm bảo mỗi thông điệp được xử lý chính xác một lần (exactly-once) là cái đích cuối cùng của tính toàn vẹn dữ liệu. Bản thân Pub/Sub theo mặc định chỉ cam kết phân phối ít nhất một lần (at-least-once delivery). Điều này có nghĩa là trong các điều kiện mạng không ổn định, lỗi phần cứng, hoặc khi xảy ra sự cố sập vùng, việc gửi trùng lặp thông điệp chắc chắn sẽ xảy ra.

Nếu không có một cơ chế xử lý exactly-once ở tầng ứng dụng hoặc tầng công cụ xử lý dữ liệu:

  • Các hệ thống kế toán, tính hóa đơn, hoặc phân tích số liệu thống kê sẽ đưa ra kết quả sai lệch (ví dụ: cộng tiền hai lần cho cùng một giao dịch).
  • Hệ thống phải gánh chịu chi phí xử lý và lưu trữ dữ liệu rác tăng vọt do các bản ghi trùng lặp liên tục được nhân bản downstream.
  • Việc tự thiết kế một cơ chế exactly-once thủ công trên ứng dụng microservice thường cực kỳ phức tạp, dễ xảy ra lỗi race condition và làm suy giảm hiệu năng hệ thống nghiêm trọng.

Sự kết hợp giữa Pub/Sub và Google Cloud Dataflow (Apache Beam) cung cấp một giải pháp out-of-the-box để đạt được exactly-once processing một cách tự động và tối ưu nhất trên đám mây GCP.


Internal Model: Cơ Chế Vận Hành Của Dataflow Exactly-Once

Để đạt được trạng thái exactly-once xử lý trên một luồng dữ liệu at-least-once từ Pub/Sub, Dataflow không thực hiện phép màu nào cả. Nó dựa trên một mô hình kiến trúc có trạng thái kết hợp chặt chẽ giữa: Message ID Tracking, State Checkpointing, và Atomic Commits.

mermaid
sequenceDiagram
    autonumber
    participant PubSub as Pub/Sub Server
    participant Worker as Dataflow Worker
    participant State as Dataflow State (Shuffle/Storage)
    participant Sink as Downstream DB (BigQuery/Spanner)

    PubSub->>Worker: Phân phối Message (ID: M1)
    Note over Worker: Đọc Message ID M1
    Worker->>State: Kiểm tra ID M1 có trong bảng Lịch Sử Khử Trùng không?
    State-->>Worker: Chưa tồn tại
    Worker->>Worker: Thực hiện tính toán biến đổi dữ liệu (Transform)
    Worker->>State: Ghi State thay đổi + Lưu ID M1 vào bảng Lịch Sử
    Worker->>Sink: Ghi kết quả downstream (Atomic/Tx)
    Worker->>State: Thực hiện Checkpoint (Commit State)
    State-->>Worker: Commit OK
    Worker->>PubSub: Gửi tín hiệu ACK cho M1
    Note over Worker: Hoàn thành xử lý M1

1. Cơ Chế Khử Trùng Lặp Ở Đầu Vào (Ingress Deduplication)

Mỗi thông điệp được xuất bản vào Pub/Sub đều được gán một thuộc tính duy nhất trên toàn cầu gọi là message_id ở tầng header của giao thức.

  • Khi Dataflow đọc dữ liệu từ Pub/Sub, nó tự động trích xuất message_id này.
  • Dataflow duy trì một bảng trạng thái lưu trữ danh sách các message_id đã được xử lý thành công.
  • Trước khi chuyển một tin nhắn qua các bước xử lý biến đổi (transformations) tiếp theo trong pipeline, worker của Dataflow sẽ truy vấn bảng trạng thái này. Nếu phát hiện message_id đã tồn tại (do Pub/Sub gửi lại vì đứt gãy mạng hoặc timeout ACK trước đó), worker sẽ bỏ qua ngay lập tức (drop duplicate) và gửi tín hiệu ACK ngược lại cho Pub/Sub để giải phóng tin nhắn trên server.

2. Trạng Thái Lưu Vết (Checkpointing) và Cam Kết Nguyên Tử (Atomic Commits)

Dataflow sử dụng cơ chế checkpointing để lưu vết trạng thái xử lý dữ liệu định kỳ hoặc theo nhóm dữ liệu (bundles).

  • Trong hệ thống streaming, dữ liệu được gom thành các bó (bundles) để xử lý song song trên nhiều workers.
  • Khi một bundle được xử lý hoàn tất, Dataflow sẽ thực hiện một giao dịch cam kết nguyên tử (atomic commit) lên hệ thống lưu trữ trạng thái (State/Shuffle layer). Giao dịch này bao gồm:
    1. Cập nhật các thay đổi trạng thái của ứng dụng (ví dụ: các biến cộng dồn, bộ đệm cửa sổ thời gian).
    2. Ghi nhận danh sách các message_id trong bundle vào bảng lịch sử khử trùng lặp.
    3. Ghi dữ liệu đầu ra downstream (nếu sink hỗ trợ ghi nhận theo giao dịch hoặc ghi đè idempotent).
  • Chỉ sau khi giao dịch commit này thành công tuyệt đối trên hệ thống lưu trữ trạng thái của Dataflow, worker mới gửi yêu cầu xác nhận (ACK) hàng loạt về cho Pub/Sub.
  • Nếu worker bị sập giữa chừng: Giao dịch commit chưa được hoàn tất. Trạng thái của Dataflow sẽ tự động rollback về checkpoint gần nhất. Pub/Sub do chưa nhận được ACK sẽ phân phối lại các message này cho worker mới. Worker mới khi nhận lại tin nhắn sẽ xử lý lại bình thường mà không sợ làm sai lệch trạng thái hệ thống, vì các thay đổi chưa được lưu vết của worker cũ đã bị hủy bỏ hoàn toàn.

Constraints, Trade-offs & Failure Modes

1. Giới Hạn Của Cửa Sổ Khử Trùng Lặp (Deduplication Window)

Đây là một ràng buộc kỹ thuật cực kỳ quan trọng mà các kỹ sư thiết kế hệ thống phải lưu ý:

  • Để tránh việc bảng lịch sử message_id phình to vô hạn và ăn mòn bộ nhớ của cụm Dataflow, Dataflow thiết lập một cửa sổ thời gian khử trùng lặp trượt (sliding deduplication window), mặc định thường là 10 phút.
  • Điều này có nghĩa là Dataflow chỉ cam kết phát hiện trùng lặp tuyệt đối nếu thông điệp bị gửi lại xuất hiện trong vòng 10 phút kể từ lần đầu tiên nó được ghi nhận.
  • Nếu một sự cố mạng nghiêm trọng xảy ra khiến Pub/Sub gửi lại một thông điệp cũ sau 11 phút, Dataflow sẽ không thể nhận diện được đây là tin nhắn trùng lặp nữa vì ID của nó đã bị xóa khỏi cache của bảng trạng thái. Pipeline sẽ xử lý thông điệp này như một tin nhắn mới, dẫn đến trùng lặp dữ liệu downstream.
  • Giải pháp: Nếu luồng dữ liệu của bạn có độ trễ lớn hoặc khả năng phục hồi sau sự cố kéo dài hơn 10 phút, bạn phải tăng cấu hình thời gian sống của bảng khử trùng lặp hoặc triển khai cơ chế khử trùng lặp ở tầng ứng dụng (Application-Level Deduplication) sử dụng một cơ sở dữ liệu có trạng thái lâu dài (như Spanner hoặc Redis) với TTL lớn hơn.

2. Sự Hạn Chế Của Các Sinks Đầu Ra (Downstream Sinks Constraint)

Tính năng exactly-once của Dataflow chỉ đảm bảo dữ liệu được xử lý chính xác một lần bên trong pipeline của nó. Khi dữ liệu đi ra ngoài (downstream write):

  • Nếu bạn ghi dữ liệu vào các hệ thống lưu trữ gốc của GCP có hỗ trợ giao tiếp transactional chính xác (như BigQuery với Write API hoặc Cloud Spanner), Dataflow sẽ phối hợp để đảm bảo exactly-once ghi đầu ra.
  • Nếu bạn sử dụng các custom sinks (ví dụ gửi HTTP POST request đến một API bên thứ ba, hoặc ghi vào một database SQL truyền thống không có tích hợp giao dịch hai pha), Dataflow không thể bảo đảm exactly-once cho các hệ thống ngoại vi này. Nếu worker bị sập ngay sau khi gọi API ngoại vi nhưng trước khi checkpoint trạng thái, worker mới sẽ chạy lại và gọi API đó lần thứ hai.
  • Quy tắc vàng: Mọi sink ngoại vi bắt buộc phải được thiết kế để xử lý dữ liệu dạng idempotent (ví dụ dùng câu lệnh UPSERT thay vì INSERT, hoặc sử dụng UUID của message làm khóa chính chống ghi trùng).

Pub/Sub Native Exactly-Once Delivery

Từ năm 2022, Google Cloud giới thiệu tính năng Exactly-Once Delivery ngay trên tầng Subscription của Pub/Sub (native feature). Đây là một tùy chọn hữu ích cho các ứng dụng microservices không sử dụng Dataflow nhưng vẫn muốn giảm thiểu trùng lặp.

yaml
# Định nghĩa Subscription với Native Exactly-Once bằng Terraform
resource "google_pubsub_subscription" "exactly_once_sub" {
  name  = "payment-processing-sub"
  topic = google_pubsub_topic.payments.name

  # Kích hoạt tính năng Exactly-Once Delivery gốc
  enable_exactly_once_delivery = true

  # Khuyến nghị tăng ACK deadline để tránh redelivery ngoài ý muốn
  ack_deadline_seconds = 60
}

Cơ chế hoạt động của Native Exactly-Once:

  1. Khóa Tin Nhắn Đang Xử Lý: Khi Pub/Sub phân phối một message cho một subscriber instance, nó sẽ khóa message đó lại. Không có subscriber nào khác có thể nhận được message này cho đến khi ACK deadline hết hạn.
  2. Xác Minh ACK ID Hợp Lệ: Nếu subscriber xử lý chậm và quá ACK deadline, hệ thống sẽ sinh ra một ACK ID mới. Khi subscriber cũ cố gắng gửi ACK bằng ID cũ, Pub/Sub sẽ từ chối thẳng thừng (ACK thất bại) vì ID đó đã bị vô hiệu hóa khi tin nhắn được giao cho instance mới.
  3. Phản Hồi Trạng Thái ACK: Khác với mặc định (gửi ACK theo kiểu fire-and-forget), khi bật exactly-once, API của Pub/Sub sẽ trả về kết quả thành công hay thất bại của lệnh ACK để ứng dụng biết chắc chắn message đã được ghi nhận hay chưa.

Giới hạn của Native Exactly-Once:

  • Chỉ hỗ trợ Pull/StreamingPull: Không hỗ trợ cho các Push Subscriptions hoặc Export Subscriptions (BigQuery/GCS export).
  • Phạm vi vùng (Regional Scope): Cam kết khử trùng lặp chỉ hoạt động tối ưu khi Publisher và Subscriber giao tiếp trong cùng một region. Nếu có định tuyến chéo vùng hoặc chuyển vùng do sự cố, khả năng trùng lặp vẫn có thể xảy ra ở tỷ lệ thấp.

References