Mô Hình Bên Trong Cloud Tasks — Task Queue Architecture
Tại Sao Phần Này Quan Trọng
Cloud Tasks là explicit task queue, không phải event stream. Sự khác biệt này ảnh hưởng mọi thứ:
- Invocation model: publisher kiểm soát chính xác khi nào, ở đâu task chạy
- Task là entity riêng, không phải event transient
- Queue là FIFO-ordered resource (có tùy chọn, nhưng default là FIFO)
- Handler phải idempotent vì task có thể execute > 1 lần
Nếu bạn không hiểu model này, bạn sẽ thiết kế sai:
- Bạn sẽ cố dùng Cloud Tasks như Pub/Sub (multiple subscribers — không được)
- Bạn sẽ bỏ qua idempotency requirement (gây duplicate data)
- Bạn sẽ miss các failure modes (task stuck in queue, never retried)
Task Queue Model vs Event Model
Pub/Sub — Implicit Invocation (Event Model)
Publisher Pub/Sub Subscribers
| | |
+----publish(message)-----------> Queue ----(implicit)-------> Handler A
| | |
| +-----------(implicit)-------> Handler B
|
(Publisher không biết Handler A, B tồn tại)Đặc điểm:
- Publisher không biết, không quan tâm subscriber
- Pub/Sub quyết định "ai xử lý message này"
- Multiple subscribers nhận cùng message
- Implicit = subscription nằm ngoài control của publisher
Cloud Tasks — Explicit Invocation (Task Queue Model)
Application Cloud Tasks Queue HTTP Handler
| | |
+--create(Task with URL)------> Queue ----(explicit)------> POST /handler
| |
| (biết chính xác nơi deliver) +---(retry if fail)-------> POST /handler
|
(Application kiểm soát destination, retry, timing)Đặc điểm:
- Application chỉ định chính xác HTTP endpoint (target)
- Task là object riêng, có lifecycle riêng
- Publisher kiểm soát retry, rate, schedule
- Single target per task (không multiple subscribers)
Khái Niệm Task & Queue
Task — Đơn Vị Công Việc
Task là struct với:
type Task struct {
// Định danh
Name string // Format: projects/X/locations/Y/queues/Z/tasks/TASKID
// HTTP Target
HttpRequest HttpRequest // Nơi deliver task
// Scheduling
ScheduleTime time.Time // Khi nào execute (default: now)
CreateTime time.Time // Khi được tạo (system)
// Dispatch Metadata
DispatchCount int32 // Bao nhiêu lần đã dispatch
ResponseCount int32 // Bao nhiêu lần nhận response
// State
Status Status // RUNNING, COMPLETED
}Lifecycle task:
[Created]
↓ (schedule_time reached)
[Dispatched]
↓ (handler returns 200-299)
[Completed + Deleted]
hoặc:
[Created]
↓ (schedule_time reached)
[Dispatched]
↓ (handler returns 5xx or timeout)
[Retry scheduled]
↓ (exponential backoff)
[Dispatched again]
↓ (max attempts exceeded)
[Dropped] (nếu không có DLQ config)Queue — Container của Tasks
Queue là resource Group tasks:
type Queue struct {
// Định danh
Name string // projects/X/locations/Y/queues/Z
State State // RUNNING, PAUSED, DISABLED
// Rate Limiting
RateLimits RateLimits // QPS, burst, concurrent dispatch
// Retry Config
RetryConfig RetryConfig // Max attempts, backoff
// Scheduling
ScheduleConfig ScheduleConfig // Timezone, max future days
}Queue là regional resource. Project có thể có 1,000 queues per region.
Dispatch Model — Làm Thế Nào Task Được Thực Thi
Step 1: Task Created
client.CreateTask(ctx, &TaskRequest{
Parent: "projects/my-project/locations/us-central1/queues/my-queue",
Task: &Task{
HttpRequest: &HttpRequest{
Url: "https://api.example.com/webhook",
Body: []byte(`{"id": "12345"}`),
Headers: map[string]string{
"Content-Type": "application/json",
"Authorization": "Bearer token-xyz",
},
},
ScheduleTime: now + 5*time.Minute, // Delay 5 min
},
})Task được lưu vào Datastore (Google's distributed database) với:
- ID (duy nhất trong queue)
- Payload, headers, URL
schedule_time= khi nào thực thiattempt_count= 0 (chưa dispatch)
Step 2: Waiting for Schedule Time
Queue service định kỳ scan tasks với schedule_time <= now. Khi schedule time tới:
Dispatch Engine kiểm tra:
- Queue state = RUNNING? (không PAUSED)
- Rate limit chưa vượt?
- Có concurrent dispatch slot trống?
Nếu tất cả OK → transition task to DISPATCHING state.
Step 3: HTTP Dispatch
Dispatch engine gửi HTTP request:
POST /webhook HTTP/1.1
Host: api.example.com
X-CloudTasks-QueueName: projects/my-project/locations/us-central1/queues/my-queue
X-CloudTasks-TaskName: projects/my-project/locations/us-central1/queues/my-queue/tasks/my-task-id
X-CloudTasks-TaskRetryCount: 0
X-CloudTasks-TaskETA: 1719360000
Content-Type: application/json
{"id": "12345"}Headers được inject tự động:
X-CloudTasks-QueueName— queue nameX-CloudTasks-TaskName— task IDX-CloudTasks-TaskRetryCount— retry attempt (0 on first)X-CloudTasks-TaskETA— when scheduled to execute (seconds since epoch)
Handler không được dùng headers này để verify identity. Đây chỉ là metadata.
Step 4: Handler Response
Success (200-299):
HTTP/1.1 200 OK→ Task được deleted (completed)
Failure (non-200-299):
HTTP/1.1 500 Internal Server Error→ Task marked for retry, exponential backoff được schedule
Timeout (> queue timeout, default App Engine 10 min): → Treat as failure → retry scheduled
Rate Limiting & Dispatch Concurrency
Queue dispatch engine không dispatch tất cả tasks một lúc. Nó kiểm soát qua 3 knobs:
1. max_dispatches_per_second (Default: 500 QPS)
Queue dispatch không vượt quá 500 tasks/second.
Ảnh hưởng:
- Nếu bạn có 10,000 tasks pending, chúng sẽ dispatch qua 20 seconds (10k / 500)
- Nếu handler xử lý lâu, tasks có thể queue up
2. max_burst_size (Default: 100)
Trong 1 second, dispatch có thể gửi tối đa 100 tasks một lúc (burst).
Ảnh hưởng:
- Nếu bạn set 500 QPS + 100 burst, second đầu dispatch 100 tasks
- Seconds tiếp theo: 100, 100, 100, ... (cân bằng để không quá burst)
3. max_concurrent_dispatches (Default: 1,000)
Có tối đa 1,000 dispatch requests flying in parallel tại bất kỳ thời điểm nào.
Ảnh hưởng:
- Nếu handler slow (10 giây/request), sau 1,000 requests, queue stalls (chờ slot)
- Nếu handler fast (100ms), 1,000 concurrent không là vấn đề
Deduplication Window — 24 Giờ Sau Deletion
Cloud Tasks tracks deleted tasks trong 24 giờ:
Task created: 2026-06-24 10:00
Task completed & deleted: 2026-06-24 10:05
Dedup window: 2026-06-24 10:05 → 2026-06-25 10:05
Nếu bạn tạo task với cùng ID trong cửa sổ này:
CreateTask(task_name="...tasks/my-task-id")
→ ALREADY_EXISTS errorTại sao? Để support idempotent creation:
// Caller muốn gửi payment task, nhưng không biết có thành công hay không
// Caller retry lần 2
for i := 0; i < 3; i++ {
err := client.CreateTask(ctx, &TaskRequest{
Task: &Task{
Name: "...tasks/payment-12345", // Fixed ID
...
},
})
if err == nil {
break
}
time.Sleep(1 * time.Second)
}Lần retry 2: CreateTask thấy task payment-12345 đã exist → error ALREADY_EXISTS (không duplicate).
Constraint quan trọng: Task ID phải duy nhất trong queue scope. Nếu bạn cố tạo task với ID đã exist (trong 24h), nó fail.
Idempotency — Bắt Buộc Implement
Cloud Tasks guarantee at-least-once delivery, không exactly-once:
Dispatch 1: Handler nhận, xử lý → database INSERT
Handler gửi response → HTTP 200
(Nhưng response bị drop trước khi Cloud Tasks nhận)
Cloud Tasks timeout → assume failed
Retry dispatch 2: Handler nhận again → database INSERT (duplicate!)
xử lý againHandler phải idempotent:
func PaymentHandler(w http.ResponseWriter, r *http.Request) {
// Trace ID từ payload
var req struct {
PaymentID string `json:"payment_id"`
Amount int `json:"amount"`
}
json.NewDecoder(r.Body).Decode(&req)
// Key: Kiểm tra xem payment đã xử lý hay không
existing, err := db.GetPayment(ctx, req.PaymentID)
if err == nil && existing.Status == "COMPLETED" {
// Duplicate request, return 200 OK (success, không re-process)
w.WriteHeader(http.StatusOK)
return
}
// First time: process
err = db.CreatePayment(ctx, req.PaymentID, req.Amount)
if err != nil {
w.WriteHeader(http.StatusInternalServerError)
return
}
// SUCCESS
w.WriteHeader(http.StatusOK)
}Idempotency key: payment_id là primary key, database reject duplicate insert. Handler return 200 OK.
Why Explicit Invocation Matters
Scenario: User clicks "Send Newsletter"
Cloud Tasks (Explicit):
// App biết target = newsletter service
for subscriber := range subscribers {
client.CreateTask(ctx, &TaskRequest{
Parent: "projects/X/locations/Y/queues/newsletter",
Task: &Task{
HttpRequest: &HttpRequest{
Url: fmt.Sprintf("https://newsletter-service.internal/send?subscriber=%s", subscriber),
},
},
})
}
// App kiểm soát: rate (500/sec), retry (10 attempts), schedule (immediately)Pub/Sub (Implicit):
// App publish "newsletter.send" event
topic.Publish(ctx, &pubsub.Message{
Data: []byte(fmt.Sprintf(`{"subscriber": "%s"}`, subscriber)),
})
// Pub/Sub kiểm soát delivery (push to subscribers)
// Subscriber có thể là CloudRun, Dataflow, Cloud Functions, ...
// App không biết ai xử lý, không kiểm soát retryKhi nào chọn nào?
| Attribute | Cloud Tasks | Pub/Sub |
|---|---|---|
| Multiple handlers | Không | Có |
| Explicit rate control | Có | Không |
| Schedule tasks | Có (30 days) | Không |
| Deduplication | Có (24h) | Không (built-in) |
| Message size | 1 MB | 10 MB |
| Ordered delivery | Có (FIFO) | Có (per key) |
Mental Model — Task State Machine
┌─────────────┐
│ CREATED │
└──────┬──────┘
│
(schedule_time reached)
│
↓
┌─────────────────┐
│ DISPATCHING │
└────┬────────┬───┘
│ │
(200-299 response)│ │(non-200 or timeout)
↓ ↓
┌─────────┐ ┌─────────────┐
│COMPLETED│ │FAILED_RETRY │
└────┬────┘ └──────┬──────┘
│ │
(delete async) │(exponential backoff)
│ │
↓ ↓
[Deleted] ┌────────────────┐
(within │ DISPATCH_AGAIN │
24 hours └────┬───────┬───┘
for dedup) │ │
│ │(max attempts?)
│ ↓
│ [Dropped/Error]
│
(retry schedule reached)
│
↓
[Dispatching]