Skip to content

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):

AttributeTypeÝ nghĩa
idStringID duy nhất của event — Eventarc generate UUID
sourceURINguồn phát sinh event, ví dụ //storage.googleapis.com/projects/my-project/buckets/my-bucket
specversionStringPhiên bản CloudEvents spec — hiện tại là "1.0"
typeStringLoại event, ví dụ google.cloud.storage.object.v1.finalized

Optional standard attributes:

AttributeTypeÝ nghĩa
timeTimestamp (RFC3339)Thời điểm event xảy ra
datacontenttypeStringMIME type của data, ví dụ application/json hay application/protobuf
subjectStringSubject của event trong context của source — ví dụ object path trong bucket
dataschemaURIURI tham chiếu đến schema của data

GCP-specific extension attributes (thêm vào cho Audit Log events):

AttributeÝ nghĩa
servicenameGCP API service name, ví dụ storage.googleapis.com
methodnameMethod trong Audit Log, ví dụ storage.objects.create
resourcenameFull 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 AuditLog message
  • 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
http
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:

python
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 '', 200

Hoặc dùng CloudEvents Python SDK để tự động parse:

python
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 '', 200

At-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:

  1. Event đến → Pub/Sub deliver cho Eventarc delivery agent
  2. Eventarc agent gọi HTTP POST đến destination
  3. Destination phải trả về 2xx response trong ack deadline
  4. 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.

python
# 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ề 2xx nhưng xử lý sai nội dung → đây là lỗi application, Eventarc không biết và không retry
  • Destination trả về 2xx như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):

bash
gcloud pubsub subscriptions modify-push-config SUBSCRIPTION_NAME \
  --push-backoff-policy=exponential \
  --push-min-retry-delay=1s \
  --push-max-retry-delay=600s

Tuy 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:

bash
# 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=5

Khi dead letter trigger:

  • Pub/Sub đã thử deliver message max-delivery-attempts lầ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ử deliver
  • CloudPubSubDeadLetterOriginalMessageId: 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:

bash
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

python
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 200

Dù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:

python
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:

python
# 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ạo

Hệ 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:

python
# 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ức

Eventarc 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.

References