CloudEvents Format & Delivery Guarantees
(File này merge CloudEvents format và delivery guarantees vì hai chủ đề này gắn chặt với nhau: format xác định cái gì được deliver, còn delivery semantics xác định cách nó được deliver và điều gì xảy ra khi gặp lỗi. Tách ra sẽ tạo ra hai file mà mỗi cái đều cần đề cập đến cái kia.)
CloudEvents: Tại sao cần standardize
Trước CloudEvents, mỗi GCP service emit event theo format riêng. Cloud Storage notifications có một schema, Pub/Sub messages có schema khác, Cloud Audit Logs có schema khác nữa. Consumer code phải biết từng format, từng field name, từng encoding rule.
CloudEvents là CNCF open standard giải quyết vấn đề này bằng cách định nghĩa một envelope chung cho mọi event: một tập context attributes mô tả event (ai, cái gì, khi nào, từ đâu), tách biệt với data payload thực sự.
Eventarc dùng CloudEvents như ngôn ngữ chung: tất cả events từ 130+ GCP providers đều được normalize về CloudEvents format trước khi deliver đến destination. Consumer chỉ cần hiểu CloudEvents, không cần biết event đến từ Cloud Storage hay BigQuery hay Audit Logs.
Cấu trúc CloudEvents
Một CloudEvents event có hai phần:
Context Attributes
Metadata mô tả sự kiện. Chia làm ba nhóm:
Required attributes (phải có trong mọi event):
| Attribute | Type | Ý nghĩa |
|---|---|---|
id | String | ID duy nhất của event — Eventarc generate UUID |
source | URI | Nguồn phát sinh event, ví dụ //storage.googleapis.com/projects/my-project/buckets/my-bucket |
specversion | String | Phiên bản CloudEvents spec — hiện tại là "1.0" |
type | String | Loại event, ví dụ google.cloud.storage.object.v1.finalized |
Optional standard attributes:
| Attribute | Type | Ý nghĩa |
|---|---|---|
time | Timestamp (RFC3339) | Thời điểm event xảy ra |
datacontenttype | String | MIME type của data, ví dụ application/json hay application/protobuf |
subject | String | Subject của event trong context của source — ví dụ object path trong bucket |
dataschema | URI | URI tham chiếu đến schema của data |
GCP-specific extension attributes (thêm vào cho Audit Log events):
| Attribute | Ý nghĩa |
|---|---|
servicename | GCP API service name, ví dụ storage.googleapis.com |
methodname | Method trong Audit Log, ví dụ storage.objects.create |
resourcename | Full resource name, ví dụ /projects/my-project/buckets/my-bucket/objects/file.csv |
Extension attributes không phải CloudEvents standard — chúng là additions của GCP để truyền thêm context cần thiết cho filter và routing.
Data Payload
Là nội dung thực sự của event. Format phụ thuộc vào provider:
- Cloud Storage events: JSON với các field
bucket,name,size,contentType, v.v. - Audit Log events: protobuf-encoded
AuditLogmessage - Pub/Sub events: wrapper với
message.data(base64-encoded),message.attributes,subscription
Provider nào emit protobuf payload sẽ có datacontenttype: application/protobuf. Consumer cần dùng proto descriptor tương ứng để decode. Đây là lý do Google Cloud Client Libraries tiện hơn tự parse — chúng đã handle proto deserialization.
HTTP Binary Mode
Eventarc deliver events qua HTTP POST sử dụng binary content mode theo CloudEvents HTTP binding spec. Điểm mấu chốt:
- Context attributes được đặt trong HTTP headers, với prefix
ce- - Data payload là HTTP body
POST /my-endpoint HTTP/1.1
Host: my-cloud-run-service.run.app
Content-Type: application/json
Authorization: Bearer eyJhbGciOiJSUzI1NiIsImtpZCI6...
ce-specversion: 1.0
ce-id: 7b6e8e97-a234-4f9c-8c43-1a2b3c4d5e6f
ce-source: //storage.googleapis.com/projects/my-project/buckets/my-bucket
ce-type: google.cloud.storage.object.v1.finalized
ce-time: 2025-06-01T12:00:00.000Z
ce-subject: objects/uploads/file.csv
ce-datacontenttype: application/json
{
"kind": "storage#object",
"id": "my-bucket/uploads/file.csv/1234567890",
"bucket": "my-bucket",
"name": "uploads/file.csv",
"size": "102400",
"contentType": "text/csv"
}Tại sao binary mode thay vì structured mode?
Structured mode đặt tất cả (cả attributes lẫn data) trong một JSON body duy nhất. Binary mode cho phép:
- HTTP infrastructure (middleware, gateway) đọc và xử lý attributes mà không cần parse body
- Consumer nhận data payload trực tiếp mà không cần unwrap
- Hiệu suất tốt hơn cho large payloads
Khi implement consumer (Cloud Run service), bạn đọc CloudEvents attributes từ HTTP headers:
from flask import request
@app.route('/', methods=['POST'])
def handle_event():
event_type = request.headers.get('ce-type')
event_id = request.headers.get('ce-id')
source = request.headers.get('ce-source')
data = request.get_json()
# Process event...
return '', 200Hoặc dùng CloudEvents Python SDK để tự động parse:
from cloudevents.http import from_http
@app.route('/', methods=['POST'])
def handle_event():
event = from_http(request.headers, request.get_data())
print(f"Type: {event['type']}, Source: {event['source']}")
return '', 200At-Least-Once Delivery: Cơ chế và hệ quả
Đây là điểm nhiều người không hiểu rõ: at-least-once delivery trong Eventarc không phải là lựa chọn thiết kế tùy ý — nó là hệ quả tất yếu của việc dùng Pub/Sub làm transport.
Pub/Sub dùng lease-based delivery model: khi một message được deliver đến consumer, consumer có một khoảng thời gian (ack deadline) để xác nhận đã xử lý xong (ack). Nếu ack deadline expire mà chưa nhận được ack → Pub/Sub re-deliver message.
Trong ngữ cảnh Eventarc:
- Event đến → Pub/Sub deliver cho Eventarc delivery agent
- Eventarc agent gọi HTTP POST đến destination
- Destination phải trả về
2xxresponse trong ack deadline - Nếu không → event được re-deliver, destination nhận lại HTTP POST
Khi nào duplicate xảy ra?
- Destination xử lý xong nhưng response bị mất do network partition → Pub/Sub re-deliver
- Destination xử lý chậm vượt ack deadline → Pub/Sub re-deliver trong khi destination vẫn đang xử lý lần đầu
- Pub/Sub re-delivers sau node failure, dù ack đã được nhận trước đó (rare nhưng có thể xảy ra)
Hệ quả: Mọi event consumer trong Eventarc phải idempotent — xử lý cùng một event nhiều lần không được tạo ra kết quả khác nhau. Dùng ce-id làm deduplication key trong state store là cách phổ biến nhất.
# Idempotent handler pattern
def handle_event(event_id, data):
# Check if already processed
if state_store.exists(event_id):
return # Already handled, skip
# Process the event
process(data)
# Mark as processed
state_store.set(event_id, processed=True)Retry Mechanics
Retry behavior của Eventarc Standard phản ánh chính xác retry behavior của Pub/Sub subscription, vì Eventarc dùng Pub/Sub subscription làm delivery mechanism.
Default retry configuration:
- Retention duration: 24 giờ — event không được ack trong 24h sẽ bị drop
- Backoff: exponential backoff
- Retry policy: tự động, không cần cấu hình thêm
Điều gì trigger retry?
- Destination trả về non-2xx HTTP status code
- Timeout (destination không trả về response trong thời gian cho phép)
- Destination không reachable (network error)
Điều gì KHÔNG trigger retry?
- Destination trả về
2xxnhưng xử lý sai nội dung → đây là lỗi application, Eventarc không biết và không retry - Destination trả về
2xxnhưng sau đó fail → Eventarc đã coi là delivered thành công
Cấu hình retry custom:
Eventarc tạo một Pub/Sub subscription ẩn cho mỗi trigger. Bạn có thể customize retry policy bằng cách sửa trực tiếp subscription này (cần tìm tên subscription trong Pub/Sub console):
gcloud pubsub subscriptions modify-push-config SUBSCRIPTION_NAME \
--push-backoff-policy=exponential \
--push-min-retry-delay=1s \
--push-max-retry-delay=600sTuy nhiên, đây là cách không được document chính thức và có thể bị override bởi Eventarc khi có thay đổi. Cách đúng đắn hơn là thiết kế destination để handle gracefully khi cần retry.
Dead Letter Handling
Eventarc Standard không có native dead letter configuration — thay vào đó, dead letter được implement thông qua Pub/Sub subscription underlying.
Cách setup dead letter cho Eventarc trigger:
# Tìm tên Pub/Sub subscription của trigger
gcloud eventarc triggers describe MY_TRIGGER --format="value(transport.pubsub.subscription)"
# Output: projects/my-project/subscriptions/eventarc-us-central1-my-trigger-sub-xxx
# Update subscription để thêm dead letter topic
gcloud pubsub subscriptions modify-push-config eventarc-us-central1-my-trigger-sub-xxx \
--dead-letter-topic=projects/my-project/topics/my-dead-letter-topic \
--max-delivery-attempts=5Khi dead letter trigger:
- Pub/Sub đã thử deliver message
max-delivery-attemptslần (5–100) - Mỗi lần thử đều nhận non-2xx response hoặc timeout
Message trong dead letter topic được wrap với additional attributes:
CloudPubSubDeadLetterSourceSubscription: subscription đã thử deliverCloudPubSubDeadLetterOriginalMessageId: ID của original message
Lưu ý về max-delivery-attempts:
Con số này là approximate. Pub/Sub không đảm bảo chính xác N lần thử trước khi forward đến DLT. Thực tế có thể ít hơn hoặc nhiều hơn một chút do distributed system timing.
IAM cho Dead Letter Topic:
Pub/Sub service account cần quyền publish vào dead letter topic:
PUBSUB_SA="service-PROJECT_NUMBER@gcp-sa-pubsub.iam.gserviceaccount.com"
gcloud pubsub topics add-iam-policy-binding my-dead-letter-topic \
--member="serviceAccount:${PUBSUB_SA}" \
--role="roles/pubsub.publisher"
# Và quyền subscribe từ source subscription để acknowledge
gcloud pubsub subscriptions add-iam-policy-binding SOURCE_SUBSCRIPTION \
--member="serviceAccount:${PUBSUB_SA}" \
--role="roles/pubsub.subscriber"Thiếu IAM trên dead letter topic là một trong những lỗi phổ biến nhất khi setup — dead letter sẽ âm thầm fail mà không có error rõ ràng.
Idempotency: Thiết kế Consumer đúng
At-least-once delivery có một hệ quả không thể né tránh: consumer phải idempotent. Đây không phải suggestion — đây là requirement bắt buộc khi dùng Eventarc.
Hai pattern idempotency phổ biến:
Pattern 1: Event ID deduplication
def process_event(event):
event_id = event['id'] # CloudEvents ce-id header
# Atomic check-and-set
if not dedup_store.set_if_absent(event_id, ttl=24*60*60):
# Already processed
return 200
# Safe to process
do_work(event['data'])
return 200Dùng ce-id làm key, lưu trong Redis hoặc Firestore với TTL 24 giờ (bằng Pub/Sub retention duration). Event nào đến sau TTL thì có thể ignore — đó là event quá cũ mà Pub/Sub đã drop.
Pattern 2: Natural idempotency trong operation
Nếu operation bản thân đã idempotent (ví dụ: set record theo primary key, không phải insert thêm record mới), thì không cần deduplication layer riêng:
def handle_file_processed(event):
file_path = event['subject'] # e.g., "objects/uploads/file.csv"
# Upsert is naturally idempotent
db.upsert(
table='processed_files',
key={'path': file_path},
values={'status': 'done', 'processed_at': event['time']}
)Không phải lúc nào cũng có thể dùng pattern 2 — nhưng khi có thể, nó đơn giản hơn nhiều và không cần distributed state.
Common mistakes
Không handle duplicate events:
# Anti-pattern: không idempotent
def process_payment(event):
amount = event['data']['amount']
db.execute("INSERT INTO transactions VALUES (?)", [amount])
# Nếu event được deliver 2 lần → 2 transactions được tạoHệ quả: mỗi lần Pub/Sub retry (network hiccup, slow response) → duplicate transaction. Trong production ở scale, điều này xảy ra thường xuyên hơn bạn nghĩ.
Trả về 2xx khi chưa process xong:
# Anti-pattern: ack sớm
def process_event(event):
asyncio.create_task(heavy_processing(event)) # fire and forget
return 200 # Trả về 2xx ngay lập tứcEventarc coi event là delivered thành công. Nếu heavy_processing fail sau đó → event bị mất, không có retry. Phải đợi processing hoàn thành trước khi trả về 2xx.
Không setup dead letter topic:
Không có DLT → events failed liên tục sẽ được retry cho đến hết 24h retention rồi bị drop. Không có record nào để investigate sau này. Trong production, DLT nên luôn được setup.