Kiến Trúc Pub/Sub & Message Lifecycle
Để hiểu delivery semantics, ordering, hay flow control — bạn cần hiểu trước Pub/Sub thực sự lưu trữ và vận chuyển message như thế nào. Phần lớn misconception về Pub/Sub xuất phát từ việc đặt nhầm mental model của Kafka hoặc RabbitMQ lên một hệ thống có thiết kế cơ bản khác.
Pub/Sub không phải Kafka — sự khác biệt kiến trúc cốt lõi
Kafka là một distributed commit log partitioned theo topic. Producer ghi vào partition, consumer đọc từ offset. Mỗi partition là một ordered sequence duy nhất, state của consumer được track bằng offset.
Pub/Sub có thiết kế khác hoàn toàn:
Không có partition, không có offset. Pub/Sub dùng mô hình lease-based delivery: thay vì ghi nhớ "consumer đã đọc đến đâu", Pub/Sub cấp phát từng message cho subscriber và chờ acknowledgment. Mỗi message có một ack deadline — nếu subscriber không ack trong thời gian đó, message được redeliver.
Storage per-subscription, không per-topic. Đây là điểm khác biệt quan trọng nhất về thiết kế:
- Kafka: message được lưu trên partition của topic, consumer tự track offset
- Pub/Sub: mỗi subscription có một "message backlog" riêng. Khi publisher gửi message tới topic, Pub/Sub tạo một bản copy (hoặc reference) của message đó cho mỗi subscription đang active. Mỗi subscription độc lập track trạng thái "đã deliver chưa, đã ack chưa" cho từng message.
Điều này có nghĩa là: nếu topic có 3 subscriptions, mỗi message phải được deliver và ack bởi cả 3 subscription trước khi Pub/Sub xóa dữ liệu gốc. Khái niệm "message consumed" trong Pub/Sub gắn với subscription, không gắn với topic.
Internal model: Distributed log sharded theo subscription
Pub/Sub documentation (Cloud Pub/Sub overview) mô tả hệ thống kết hợp "horizontal scalability của Apache Kafka và Pulsar với features của messaging middleware" — nhưng không đi sâu vào implementation detail. Từ behavior quan sát được và GCP architecture blog, ta có thể suy ra model sau:
Message storage layer
Bên dưới, Pub/Sub dùng Colossus (distributed filesystem của Google) để lưu message durably. Khi publisher gửi message:
- Message được ghi vào Colossus với multi-zone replication
- Chỉ sau khi write đã được acknowledge bởi ít nhất 2 availability zones trong region, Pub/Sub mới trả về
publish ackcho publisher - Message ID được assign — unique identifier dùng để track delivery state
Đây là lý do tại sao publish latency của Pub/Sub cao hơn một chút so với in-memory queue: đổi lại là durability guarantee — message không thể bị mất sau khi publisher nhận ack, kể cả khi một zone fail.
Delivery state per subscription
Với mỗi subscription, Pub/Sub maintain một delivery state cho từng message chưa được ack:
- Message ID
- Delivery attempt count
- Current ack deadline (timestamp)
- Subscriber ID đang giữ lease (nếu đang được deliver)
State này cũng được lưu durably — cho phép Pub/Sub redeliver message đúng cách ngay cả sau khi Pub/Sub server restart.
Sharding
Vì Pub/Sub phục vụ hàng triệu subscriptions và billions of messages mỗi ngày, mọi thứ phải được sharded. Sharding xảy ra ở hai tầng:
Topic sharding: Message của cùng một topic được phân phối trên nhiều "message storage server". Không có khái niệm partition order giữa các shard — đây là một trong các lý do tại sao Pub/Sub không guarantee global ordering theo mặc định.
Subscription sharding: Delivery state của một subscription có thể được phân phối trên nhiều "subscription server". Khi subscriber pull message, có thể nhận message từ nhiều server khác nhau. Đây là tại sao với ordering keys, Pub/Sub phải route messages với cùng key về cùng một server — để maintain per-key order.
Message lifecycle: Từ publish đến xóa
Vòng đời hoàn chỉnh của một message trải qua 5 giai đoạn:
Giai đoạn 1: Published
Publisher → Pub/Sub API → "Write to Colossus (multi-zone)" → Publish ACKPublisher gọi publish() với message data và optional attributes. Pub/Sub:
- Nhận message tại frontend server
- Ghi message vào storage với multi-zone replication
- Fanout message sang delivery state của tất cả subscriptions đang active
- Trả về message ID cho publisher
Quan trọng: Nếu publisher không nhận được publish ACK (timeout, network error), publisher không biết message có được lưu chưa. Pub/Sub recommend retry với backoff — điều này có thể dẫn đến duplicate publishes. Đây là một trong các nguồn của at-least-once semantics.
Giai đoạn 2: Stored
Message ở trong backlog của subscription, chưa được deliver. Trạng thái này kéo dài đến khi:
- Subscriber pull message hoặc push subscription trigger
- Message retention period hết hạn (default 7 ngày, max 31 ngày)
Message retention period là property của topic, không phải subscription. Nếu subscription không tồn tại khi message được publish, message sẽ không được deliver — không có backfill.
Giai đoạn 3: Delivered (Leased)
Subscriber ← Pub/Sub delivers message (with ack deadline)Khi subscriber pull (hoặc push endpoint nhận được HTTP request), message được lease cho subscriber đó:
- Message vẫn tồn tại trong storage
- Ack deadline bắt đầu đếm ngược
- Nếu ack deadline hết mà chưa có ack → message được redeliver
"Lease" không có nghĩa là chỉ một subscriber nhận được message. Nếu subscriber nhận message nhưng không ack và ack deadline expire, một subscriber khác (hoặc cùng subscriber đó sau khi reconnect) sẽ nhận lại message. Đây là cơ chế sinh ra duplicates trong at-least-once delivery.
Giai đoạn 4: Acknowledged
Subscriber → ack(messageId) → Pub/Sub marks message as "delivered" for this subscriptionSubscriber gửi ACK với ack ID. Pub/Sub:
- Đánh dấu message là "acknowledged" trong delivery state của subscription
- Message không còn được deliver lại cho subscription này
Lưu ý: ACK là per-subscription. Nếu topic có 3 subscriptions, message chỉ bị xóa khỏi backlog của subscription nào đã ack, không ảnh hưởng đến subscription khác.
Giai đoạn 5: Deleted
Message bị xóa khỏi storage khi tất cả subscriptions đã ack message đó, hoặc khi retention period hết hạn. Pub/Sub tự động garbage collect.
Một case đặc biệt: nếu một subscription bị xóa sau khi message được publish nhưng trước khi message được ack, message sẽ được xóa khỏi backlog của subscription đó ngay lập tức. Nếu đó là subscription duy nhất, message bị xóa sớm.
Tại sao at-least-once là default — phân tích từ góc độ distributed systems
At-least-once delivery không phải là "lỗi thiết kế" — nó là hệ quả tất yếu của việc ưu tiên availability và scalability trong một distributed system.
Để đảm bảo exactly-once, hệ thống cần:
- Biết chính xác message đã được deliver hay chưa — đòi hỏi distributed lock hoặc consensus protocol
- Biết subscriber đã xử lý thành công hay chưa — đòi hỏi hai-phase commit hoặc idempotent processing
- Maintain state này durably và consistent — overhead rất lớn
Trong mô hình at-least-once, Pub/Sub chỉ cần đảm bảo: "nếu subscriber ack, message không bị redeliver". Còn nếu ack bị mất (network partition giữa subscriber và Pub/Sub server), message sẽ được redeliver — an toàn hơn là mất message.
Lease-based model giải thích điều này rõ ràng nhất:
Timeline:
T=0s: Pub/Sub delivers message, sets ack_deadline=30s
T=10s: Subscriber xử lý xong, gửi ACK
T=11s: ACK packet bị drop do network issue
T=30s: Ack deadline expire, Pub/Sub thấy "chưa có ACK"
T=30s: Pub/Sub redeliver message → DUPLICATESubscriber đã xử lý thành công, nhưng vẫn nhận được duplicate. Đây là at-least-once: message được deliver ít nhất một lần, có thể nhiều hơn.
Để xử lý correctly, subscriber phải implement idempotent processing: xử lý cùng một message nhiều lần phải cho kết quả như xử lý một lần.
Message format và limits
Một message trong Pub/Sub có cấu trúc:
- data: byte array (raw content), base64-encoded khi dùng REST API
- attributes: key-value pairs string (metadata, không phải content)
- messageId: assigned by Pub/Sub sau khi publish
- publishTime: timestamp khi message được publish
- orderingKey: optional, string tối đa 1 KB
Limits quan trọng:
- Message size tối đa: 10 MB (data + attributes gộp lại)
- Batch publish tối đa: 10 MB hoặc 1000 messages, whichever comes first
- Attribute key/value: tối đa 256 bytes each, tối đa 100 attributes per message
Multi-zone replication và durability
Pub/Sub đảm bảo message được replicate trên ít nhất 2 zones trong region trước khi trả publish ACK. Điều này có nghĩa:
- Zone failure không gây mất message đã được acknowledged
- Nhưng region failure có thể gây gián đoạn delivery, không phải mất message
Pub/Sub là regional service: mỗi topic/subscription được deploy trong một region. Cross-region replication không tự động — nếu bạn cần multi-region durability, phải thiết kế architecture phù hợp (ví dụ: publish vào nhiều regional topics, hoặc dùng message export).
Implications cho thiết kế hệ thống
Hiểu được internal model, ta rút ra các implications thực tế:
Fanout tự động nhưng có chi phí: Nếu topic có 10 subscriptions và nhận 1M messages/ngày, storage phải chứa trạng thái cho 10M messages. Đây là lý do Pub/Sub charge theo message delivery, không chỉ theo message publish — mỗi subscription là một lần delivery.
Backlog là per-subscription: Khi bạn monitor "message backlog", phải monitor per-subscription, không per-topic. Một subscription slow consumer làm tăng backlog của nó mà không ảnh hưởng subscription khác.
Message retention là safety net, không phải replay log: Không như Kafka nơi bạn có thể seek lại offset để replay, Pub/Sub không có API seek theo timestamp/offset theo cách Kafka làm. Pub/Sub có Seek API cho phép reset subscription tới một snapshot hoặc timestamp, nhưng cơ chế này khác về bản chất và không phải cho streaming replay.
Publish failures cần retry idempotent: Vì publisher không biết message có được lưu hay không khi mạng lỗi, retry là bắt buộc. Pub/Sub message ID là unique per successful publish — nếu cùng một message được publish hai lần (do retry), sẽ có hai message IDs khác nhau và subscriber sẽ nhận cả hai. Pub/Sub không dedup ở publish layer (chỉ exactly-once delivery dedup ở delivery layer).