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=3600sQueue 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-central1Incident 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-central1Pro: 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 QPSBenefits
✓ Isolation: One shard overloaded doesn't affect others
✓ Scale: Add more shards as needed
✓ Monitoring: Per-shard metrics
✓ Failover: Can pause one shard for maintenanceMaintenance 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 anySteps
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-v1Rate 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 statsIncrease 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, revertDecrease 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 stableObservability 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 insufficientDashboard 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