Google Cloud Tasks for HTTP Workers
Cloud Tasks inverts the usual worker model: instead of workers pulling from a queue, the service pushes each task to an HTTP endpoint at a rate you configure. This guide builds a production setup with it as part of Managed Cloud Queues for Background Jobs in Backend Frameworks & Worker Scaling. It fits especially well when the hard problem is pacing — calling a rate-limited API or a fragile downstream — because the rate limit is enforced before your code ever runs.
Problem Statement
A marketplace syncs product listings to three partner channels whenever a seller edits a product. Each partner enforces its own rate limit (5, 20, and 50 requests per second) and returns 429 when exceeded. The current implementation calls partners synchronously from the API and falls over during bulk edits, when a seller updates 10,000 products at once. You want sync calls to happen in the background, each partner paced at its own limit, retries with backoff for 429 and 5xx responses, a single sync per product edit even if the API retries the enqueue, and a record of any sync that ultimately fails.
Prerequisites
- A Google Cloud project with the Cloud Tasks API enabled and a region chosen (queues are regional).
- An HTTP handler deployed somewhere reachable — Cloud Run is assumed here, with authentication required.
- A service account for Cloud Tasks to mint OIDC tokens as (
tasks-invoker@...), grantedroles/run.invokeron the handler service. - The
google-cloud-tasksclient library in the enqueueing service.
Step 1 — Create One Queue per Rate Limit
Rate limits and retry policies are properties of a queue, so partners with different limits get different queues.
for partner in acme:5 globex:20 initech:50; do
name=${partner%%:*}; rate=${partner##*:}
gcloud tasks queues create "sync-$name" \
--location=europe-west1 \
--max-dispatches-per-second="$rate" \ # the partner's limit
--max-concurrent-dispatches=$((rate * 4)) \ # allows for ~4 s handler latency
--max-attempts=15 \ # then the task is deleted
--min-backoff=10s --max-backoff=30m \
--max-doublings=8 \
--max-retry-duration=24h # stop retrying after a day regardless
done
max-dispatches-per-second is a token bucket: the queue dispatches at most that many tasks per second, with a small burst allowance. max-concurrent-dispatches caps tasks in flight; with a handler that takes up to four seconds, a concurrency of four times the rate keeps the pipe full without exceeding the rate. The retry settings define exponential backoff from 10 seconds doubling up to 30 minutes — the same shape described in exponential backoff with jitter in Celery, enforced by the service.
Step 2 — Enqueue Named Tasks from the API
Give each task a deterministic name derived from what it represents. Creating a task whose name already exists — or existed recently — fails with ALREADY_EXISTS, which turns a retried enqueue into a no-op.
# enqueue.py
import json
from google.api_core.exceptions import AlreadyExists
from google.cloud import tasks_v2
client = tasks_v2.CloudTasksClient()
HANDLER = "https://listing-sync-abc123-ew.a.run.app"
INVOKER = "tasks-invoker@my-project.iam.gserviceaccount.com"
def enqueue_sync(partner: str, product_id: str, revision: int) -> None:
parent = client.queue_path("my-project", "europe-west1", f"sync-{partner}")
task = {
"name": f"{parent}/tasks/{partner}-{product_id}-r{revision}", # dedup key
"http_request": {
"http_method": tasks_v2.HttpMethod.POST,
"url": f"{HANDLER}/sync/{partner}",
"headers": {"Content-Type": "application/json"},
"body": json.dumps({"product_id": product_id, "revision": revision}).encode(),
"oidc_token": {"service_account_email": INVOKER, "audience": HANDLER},
},
"dispatch_deadline": {"seconds": 30},
}
try:
client.create_task(parent=parent, task=task)
except AlreadyExists:
pass # a retry of the same enqueue: fine
Including the product revision in the name matters: a second edit to the same product is a new sync that must not be deduplicated away. Task names are reserved for up to about an hour after a task completes or is deleted, and sequential or hash-free names can create hotspots in the service's internal storage, so prefer names with a high-entropy prefix at very high rates.
Step 3 — Write an Authenticated, Idempotent Handler
The handler receives the task body as an ordinary HTTP request. With Cloud Run authentication required, only requests carrying a valid OIDC token for an identity with run.invoker reach your code. Cloud Tasks adds headers with the queue name, task name, and retry count.
# app.py — FastAPI on Cloud Run
from fastapi import FastAPI, Request, Response
app = FastAPI()
@app.post("/sync/{partner}")
async def sync(partner: str, request: Request) -> Response:
body = await request.json()
retry = int(request.headers.get("X-CloudTasks-TaskRetryCount", "0"))
task_name = request.headers.get("X-CloudTasks-TaskName", "")
listing = await load_listing(body["product_id"])
if listing.revision > body["revision"]:
return Response(status_code=200) # stale task: a newer sync supersedes it
status = await partners[partner].upsert(listing) # returns the partner's HTTP status
if status in (429, 500, 502, 503, 504):
return Response(status_code=503) # non-2xx: Cloud Tasks retries with backoff
if 400 <= status < 500:
await record_permanent_failure(partner, body, status, task_name)
return Response(status_code=200) # don't retry a request the partner rejects
return Response(status_code=200)
Two rules make this correct. Return 2xx only after the work is done — Cloud Tasks deletes the task on 2xx, so returning early and working in the background loses work on a crash. And return 2xx for permanent failures after recording them, because Cloud Tasks has no dead-letter queue: after the final attempt the task is simply deleted, and your own record is the only trace it existed.
Step 4 — Respect the Dispatch Deadline
dispatch_deadline is how long Cloud Tasks waits for a response (15 seconds to 30 minutes for HTTP targets; the default is 10 minutes). If the handler has not responded by then, the attempt counts as failed and is retried — even if the handler is still running. Set the deadline above the handler's worst case, and make sure the platform's own request timeout is at least as long.
# Cloud Run request timeout must be >= the task dispatch deadline
gcloud run services update listing-sync --region=europe-west1 --timeout=60s
The failure this prevents is subtle because nothing errors: a handler that takes 45 seconds under a 30-second deadline finishes its work successfully, but Cloud Tasks has already recorded a failed attempt and scheduled a retry. The retry performs the same sync again. With an idempotent handler that costs only partner quota; with a non-idempotent one it duplicates side effects on every slow request.
For work that genuinely takes longer than a comfortable deadline, split it: the handler does a bounded chunk and enqueues a follow-up task for the next chunk.
Step 5 — Schedule Tasks for Later
Tasks can carry a schedule_time up to 30 days in the future, which replaces a separate scheduler for delayed work such as "sync again in an hour if the partner was unavailable".
from google.protobuf import timestamp_pb2
import datetime as dt
def enqueue_later(partner: str, product_id: str, revision: int, delay: dt.timedelta) -> None:
ts = timestamp_pb2.Timestamp()
ts.FromDatetime(dt.datetime.now(dt.timezone.utc) + delay)
task = build_task(partner, product_id, revision) # as in Step 2
task["schedule_time"] = ts
client.create_task(parent=queue_path(partner), task=task)
Scheduled tasks count toward the queue's size but not its dispatch rate until they are due. For recurring schedules, use Cloud Scheduler to create tasks or call the handler directly; Cloud Tasks has no cron. The general scheduling trade-offs are covered in Scheduled & Delayed Jobs.
Step 6 — Pause, Purge, and Observe Queues
Queues can be paused (tasks accumulate, nothing is dispatched) and resumed — the fastest way to stop hammering a partner during an incident without losing work.
gcloud tasks queues pause sync-acme --location=europe-west1 # incident: stop dispatching
gcloud tasks queues resume sync-acme --location=europe-west1 # dispatch resumes at the set rate
gcloud tasks queues describe sync-acme --location=europe-west1 --format="value(state,rateLimits)"
In Cloud Monitoring, watch cloudtasks.googleapis.com/queue/depth and queue/task_attempt_count grouped by response_code. A rising share of 503s means the partner is throttling below its advertised limit — lower max-dispatches-per-second rather than letting retries do the pacing.
Verification
# Enqueue the same task twice: the second create must fail with ALREADY_EXISTS
python -c 'from enqueue import enqueue_sync; enqueue_sync("acme","p-1",7); enqueue_sync("acme","p-1",7)'
# Load test: 2,000 tasks into sync-acme should take ~400 s at 5/s
python bulk.py --partner acme --count 2000
gcloud logging read 'resource.type="cloud_run_revision" AND httpRequest.requestUrl:"/sync/acme"' \
--freshness=10m --format="value(timestamp)" | sort | uniq -c | awk '{print $1}' | sort -n | tail -1
# expect the max per-second count to be ~5
Gotchas & Edge Cases
No dead-letter queue. After the last attempt, the task is gone. Log at the final attempt (compare X-CloudTasks-TaskRetryCount with your max attempts) or record permanent failures in your own table, and alert on it.
At-least-once delivery. Tasks can occasionally be dispatched more than once even without failures. Keep handlers idempotent, as the revision check does.
Enqueue throughput limits. create_task has per-queue quotas; a bulk operation creating hundreds of thousands of tasks should enqueue in parallel from several workers or use a fan-out task that creates the rest.
Regional queues. Queues live in one region. A regional outage stops dispatch; for critical flows, consider the blast radius when choosing the region.
FAQ
Cloud Tasks or Pub/Sub for background jobs? Cloud Tasks when you need per-queue rate limits, explicit scheduling, and task-level deduplication for work aimed at one handler. Pub/Sub when one event feeds several independent consumers. Cloud Pub/Sub push vs pull for workers compares the Pub/Sub options.
Can Cloud Tasks call targets outside Google Cloud? Yes — HTTP targets can be any public HTTPS endpoint. Use OAuth or OIDC tokens only for Google-hosted targets; for external ones, add your own shared-secret header and verify it.
How do I prioritise some tasks over others? Use separate queues with different rates; Cloud Tasks has no per-task priority within a queue.
Related
- Managed Cloud Queues for Background Jobs — push vs pull and provider trade-offs.
- Rate Limiting Third-Party API Calls from Workers — the same pacing problem solved in the worker.
- Reliable Webhook Delivery with a Job Queue — outbound HTTP delivery with retries.
- Cloud Pub/Sub Push vs Pull for Workers — Google's other queueing service.