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@...), granted roles/run.invoker on the handler service.
  • The google-cloud-tasks client 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.

One handler, three paces A seller's bulk edit enqueues 10,000 tasks into each of three partner queues. The acme queue dispatches at 5 per second, globex at 20, and initech at 50. All three call the same Cloud Run handler, which never sees more than the combined configured rate, so no partner receives more than its limit. Rate limits enforced before the handler runs bulk edit 3 x 10,000 tasks sync-acme: 5/s sync-globex: 20/s sync-initech: 50/s Cloud Run handler POST /sync/{partner} acme drains in ~33 minutes, initech in ~3; neither sees a 429 from pacing.

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.

Status code is the contract A 2xx response tells Cloud Tasks the task succeeded and it is deleted. A non-2xx response, such as 503 returned when the partner answered 429 or 5xx, schedules a retry with exponential backoff. For a permanent partner rejection, the handler records the failure in its own database and returns 200, because Cloud Tasks deletes tasks after the last attempt without a dead-letter queue. What each response tells Cloud Tasks 200 after work done task deleted success 503 (partner 429/5xx) retry with backoff up to max attempts record, then 200 permanent 4xx from partner your table is the DLQ Returning 200 before the work finishes loses the task on any crash.

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.

Deadline shorter than the work means a duplicate Cloud Tasks dispatches a task with a 30-second deadline. The handler takes 45 seconds and succeeds. At 30 seconds Cloud Tasks gave up waiting, counted the attempt as failed, and scheduled a retry, which runs the same sync a second time. Setting the deadline above the handler's worst case prevents this. dispatch_deadline 30 s, handler 45 s attempt 1: runs 45 s, succeeds 30 s: Cloud Tasks gives up waiting attempt 2 after backoff: same sync again Deadline must exceed the handler's worst case, and Cloud Run's timeout must match.

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