Skip to content

Operational Patterns

Tại Sao Quan Trọng

Queue lifetime không chỉ "create once, forget". Production systems cần:

  • Graceful shutdown — drain pending tasks, don't lose work
  • Incident response — pause queue during outage, resume after fix
  • Maintenance — purge old tasks, rotate queues
  • Capacity management — create/delete queues, adjust limits

Queue Lifecycle

Create Queue

bash
gcloud tasks queues create my-queue \
    --location=us-central1 \
    --max-dispatches-per-second=500 \
    --max-burst-size=100 \
    --max-concurrent-dispatches=1000 \
    --max-attempts=100 \
    --min-backoff=100ms \
    --max-backoff=3600s

Queue States

RUNNING (default)
  → Queue dispatching tasks normally
  → CreateTask, DispatchTask work

PAUSED
  → Queue not dispatching
  → CreateTask still works (tasks enqueued)
  → Tasks held until resume
  → Useful for incident response

DISABLED
  → Queue deleted (deprecated state)

Pause/Resume

bash
# Pause queue (stop dispatching)
gcloud tasks queues pause my-queue --location=us-central1

# Resume queue (start dispatching)
gcloud tasks queues resume my-queue --location=us-central1

Incident Response Pattern

Scenario: Handler Bug, Cascade Failure

Time T0: Handler bug introduced
  → Tasks return 500 errors
  → Retry storm begins
  → CPU spike, cascading failures

Time T0+5min: Alert fires
  → On-call engineer woken up
  
Step 1: Pause queue (STOP the bleeding)
  gcloud tasks queues pause my-queue

  Result: No more retries, no more dispatch
  Pending tasks held in queue (safe)

Step 2: Investigate
  - Check application logs
  - Identify root cause (bad SQL query, infinite loop, etc.)
  
Step 3: Deploy fix
  - Fix code, redeploy, health check passes
  
Step 4: Resume queue
  gcloud tasks queues resume my-queue
  
  Result: Tasks resume dispatching (retry scheduled tasks)
  
Step 5: Monitor recovery
  - Check failure rate (should drop to 0%)
  - Check task age (backlog should drain)
  - Check handler latency (should normalize)

Pause Without Data Loss

go
// Application shutdown
func gracefulShutdown() {
    // Step 1: Stop accepting new task creates
    pauseTaskCreation = true
    
    // Step 2: Pause queue (don't dispatch existing tasks)
    client.UpdateQueue(ctx, &Queue{
        Name: "projects/X/locations/Y/queues/my-queue",
        State: PAUSED,
    })
    
    // Step 3: Let in-flight requests finish
    time.Sleep(30 * time.Second)
    
    // Step 4: Shutdown gracefully
    server.Shutdown(ctx)
    
    // On restart:
    // - Queue still PAUSED, tasks safe
    // - Resume when ready: gcloud tasks queues resume ...
}

Purge Pattern

Purge All Tasks from Queue

bash
# Delete and recreate queue (nuclear option)
gcloud tasks queues delete my-queue --location=us-central1 -q
gcloud tasks queues create my-queue --location=us-central1

Pro: Cleans up all tasks
Con: Loses all queued work, need to recreate queue config

Selective Purge (Delete Specific Tasks)

bash
# List and delete tasks matching condition
gcloud tasks list-tasks --queue=my-queue --location=us-central1 \
  --filter="createTime<2026-01-01" \
  --format=json | \
  jq -r '.[].name' | \
  xargs -I {} gcloud tasks delete {}

Bulk Delete Pattern

go
func purgeQueueBefore(ctx context.Context, queueName string, before time.Time) error {
    client := cloudtasks.NewClient(ctx)
    defer client.Close()
    
    it := client.ListTasks(ctx, &taskspb.ListTasksRequest{
        Parent: queueName,
    })
    
    for {
        task, err := it.Next()
        if err == iterator.Done {
            break
        }
        if err != nil {
            return err
        }
        
        if task.CreateTime.AsTime().Before(before) {
            client.DeleteTask(ctx, &taskspb.DeleteTaskRequest{
                Name: task.Name,
            })
        }
    }
    
    return nil
}

Scaling Pattern: Multi-Queue

Problem: Single Queue Bottleneck

max_dispatches_per_second = 500 (per queue limit)
But need 10,000 QPS

Solution: Create 20 queues (500 * 20 = 10,000 QPS total)

Implementation

go
// Distribute tasks across queues based on shard key
func CreateTaskSharded(ctx context.Context, req PaymentRequest) error {
    // Shard key: user ID
    shardID := hashUserID(req.UserID) % NUM_SHARDS
    
    queueName := fmt.Sprintf(
        "projects/X/locations/Y/queues/payment-shard-%d",
        shardID,
    )
    
    client.CreateTask(ctx, &TaskRequest{
        Parent: queueName,
        Task: &Task{
            HttpRequest: &HttpRequest{
                Url: "https://payment-service.internal/charge",
                Body: []byte(fmt.Sprintf(`{"payment_id": "%s"}`, req.PaymentID)),
            },
        },
    })
    
    return nil
}

// Each queue independent:
// - Queue 0: 500 QPS
// - Queue 1: 500 QPS
// - ...
// - Queue 19: 500 QPS
// - Total: 10,000 QPS

Benefits

✓ Isolation: One shard overloaded doesn't affect others
✓ Scale: Add more shards as needed
✓ Monitoring: Per-shard metrics
✓ Failover: Can pause one shard for maintenance

Maintenance Pattern: Queue Rotation

Scenario: Zero-Downtime Queue Migration

Old queue: "payment-v1" (out of date config)
New queue: "payment-v2" (updated config, faster handler)

Goal: Migrate all pending tasks without losing any

Steps

Step 1: Create new queue with updated config
  gcloud tasks queues create payment-v2 \
    --location=us-central1 \
    --max-dispatches-per-second=1000

Step 2: Enable dual-write (application creates in both queues)
  if enableDualWrite {
    createTask(queueName="payment-v1", ...)
    createTask(queueName="payment-v2", ...)
  }

Step 3: Drain old queue (let existing tasks complete)
  - Monitor payment-v1 task count
  - Wait for all tasks to dispatch and complete
  
Step 4: Disable writes to old queue
  enableDualWrite = false
  // New tasks only to payment-v2
  
Step 5: Verify old queue empty
  gcloud tasks list-tasks --queue=payment-v1
  // Should be empty
  
Step 6: Delete old queue
  gcloud tasks queues delete payment-v1

Rate Adjustment Pattern

Monitor and Adjust

bash
# Check current queue stats
gcloud tasks queues describe my-queue --location=us-central1

# Shows:
# - Current rate limits
# - Task count
# - Recent dispatch stats

Increase Rate (High Demand)

bash
# Gradually increase rate limit
gcloud tasks queues update my-queue \
    --location=us-central1 \
    --max-dispatches-per-second=1000

# Monitor handler health (latency, errors)
# If stable, increase more
# If errors spike, revert

Decrease Rate (Handler Overload)

bash
# Emergency: reduce rate
gcloud tasks queues update my-queue \
    --location=us-central1 \
    --max-dispatches-per-second=100 \
    --max-concurrent-dispatches=50

# Result: Tasks dispatch slower, handler recovers
# Backlog will grow initially, but queue becomes stable

Observability Pattern: Queue Monitoring

Key Metrics to Track

Queue Depth
  = number of tasks pending
  High depth = backlog growing
  Should normalize after incident resolved

Dispatch Rate (actual)
  = tasks dispatched per second
  Should match config (or be limited by handler)

Error Rate
  = % tasks returning non-200 status
  Baseline (normal) vs spike (incident)
  Alert if > threshold

Handler Latency
  = time from dispatch to response
  If increases, effective QPS decreases (due to concurrent limit)

Task Age
  = age of oldest pending task
  Should be low (< 1 second)
  High age = queue stuck or rate insufficient

Dashboard Setup

yaml
# Example Cloud Monitoring dashboard config
dashboard:
  - title: "Cloud Tasks Queue Health"
    widgets:
      - chart:
          title: "Queue Depth"
          metric: "cloudtasks.googleapis.com/queue/task_count"
          
      - chart:
          title: "Dispatch Rate (QPS)"
          metric: "cloudtasks.googleapis.com/queue/task_attempt_count"
          
      - chart:
          title: "Error Rate (%)"
          metric: "cloudtasks.googleapis.com/queue/task_attempt_count"
          filter: 'metric.status != "OK"'
          
      - chart:
          title: "Handler Latency (ms)"
          metric: "custom.googleapis.com/handler/latency"

Testing Pattern: Load Testing Queue

Simulate Production Load

go
func LoadTest(ctx context.Context, queueName string, rps int) {
    ticker := time.NewTicker(time.Second / time.Duration(rps))
    defer ticker.Stop()
    
    client := cloudtasks.NewClient(ctx)
    defer client.Close()
    
    for i := 0; i < 10000; i++ {
        <-ticker.C
        
        client.CreateTask(ctx, &TaskRequest{
            Parent: queueName,
            Task: &Task{
                Name: fmt.Sprintf("...tasks/load-test-%d", i),
                HttpRequest: &HttpRequest{
                    Url: "https://handler.example.com/test",
                    Body: []byte(fmt.Sprintf(`{"test_id": %d}`, i)),
                },
            },
        })
    }
}

// Usage:
// - Create test queue with production config
// - Run: LoadTest(ctx, "projects/X/queues/test", 5000)
// - Monitor handler latency, CPU, error rate
// - Measure effective throughput
// - Adjust config based on results

References