Skip to content

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:

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

go
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

go
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 thi
  • attempt_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:

http
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 name
  • X-CloudTasks-TaskName — task ID
  • X-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
HTTP/1.1 200 OK

→ Task được deleted (completed)

Failure (non-200-299):

http
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 error

Tại sao? Để support idempotent creation:

go
// 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ý again

Handler phải idempotent:

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

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

go
// 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 retry

Khi nào chọn nào?

AttributeCloud TasksPub/Sub
Multiple handlersKhông
Explicit rate controlKhông
Schedule tasksCó (30 days)Không
DeduplicationCó (24h)Không (built-in)
Message size1 MB10 MB
Ordered deliveryCó (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]

References