Celery Chains, Groups, and Chords

Celery's canvas primitives are the fastest way to express a multi-step pipeline in Python, and this guide shows how to use them safely as part of Job Chaining & Workflow Orchestration within Queue Fundamentals & Architecture. The primitives are simple to write and easy to get subtly wrong: argument passing, result backend dependencies, and error propagation in chords all behave differently from what the syntax suggests.

Problem Statement

A reporting service generates a monthly statement per customer: fetch usage from three regional databases in parallel, merge the results, render a PDF, and email it. A first attempt enqueued each step from inside the previous task, and failures left customers with no statement and no alert. You want the whole pipeline expressed as one Celery workflow, the parallel fetch joined exactly once, a failure anywhere routed to a single error handler, and a way to run it for 40,000 customers without flooding the broker.

Prerequisites

  • Celery 5.3+ with a broker (Redis or RabbitMQ) configured as in setting up Celery with a Redis broker.
  • A result backend — chords do not work without one. Redis is the common choice; the database backend works but is slower for chord counters.
  • Tasks that are idempotent, because a worker crash with acks_late redelivers a task that may have partly run.
  • result_expires long enough to outlive your slowest workflow (the default is one day).

Step 1 — Chain Sequential Steps with Immutable Signatures

A chain runs tasks one after another, passing each task's return value as the first argument of the next. That implicit passing is the most common source of confusion: if the next task does not expect the previous result, use an immutable signature (.si()) so nothing is prepended.

# tasks.py
from celery import Celery, chain

app = Celery("statements", broker="redis://redis:6379/0", backend="redis://redis:6379/1")

@app.task(bind=True, acks_late=True, max_retries=3)
def merge_usage(self, parts: list[dict], customer_id: str) -> str:
    # parts is the list of results from the preceding group (see Step 2)
    key = storage.put_json(f"usage/{customer_id}.json", combine(parts))
    return key                                    # pass a storage key, not the data

@app.task(acks_late=True)
def render_pdf(usage_key: str, customer_id: str) -> str:
    return storage.put_bytes(f"pdf/{customer_id}.pdf", render(storage.get_json(usage_key)))

@app.task(acks_late=True)
def email_statement(pdf_key: str, customer_id: str) -> None:
    mailer.send_statement(customer_id, pdf_key)

# render_pdf receives merge_usage's return value as its first argument:
tail = chain(render_pdf.s(customer_id="c-17"), email_statement.s(customer_id="c-17"))

.s() builds a mutable signature that accepts the parent's result; .si() builds one that ignores it. Mixing them up produces TypeError: takes 2 positional arguments but 3 were given at runtime, not at definition time — test workflows end to end, not just individual tasks.

What a chain passes along Task A returns a storage key that becomes the first argument of task B, built with a mutable signature. Task B returns a PDF key passed to task C. A task built with an immutable signature receives only the arguments given when the signature was created. .s() receives the parent result; .si() does not merge_usage returns usage_key render_pdf.s(...) (usage_key, customer_id) email.s(...) (pdf_key, customer_id) cleanup.si(customer_id): runs after, ignores what came before

Step 2 — Fan Out with a Group, Join with a Chord

A group runs tasks in parallel. A chord is a group plus a callback (the "body") that runs once, with the list of all group results, after every member has finished.

from celery import chord, group

@app.task(acks_late=True, autoretry_for=(DatabaseTimeout,), retry_backoff=True, max_retries=5)
def fetch_usage(region: str, customer_id: str, month: str) -> dict:
    return regional_db(region).usage(customer_id, month)   # small dict, safe to return

def statement_workflow(customer_id: str, month: str):
    fetch_all = group(fetch_usage.si(r, customer_id, month) for r in ("eu", "us", "ap"))
    return chain(
        chord(fetch_all, merge_usage.s(customer_id=customer_id)),   # body gets [eu, us, ap]
        render_pdf.s(customer_id=customer_id),
        email_statement.s(customer_id=customer_id),
    )

Under the hood, the chord stores a counter in the result backend. Each member's completion increments it; when the count reaches the group size, a worker enqueues the body. With the Redis backend this is an atomic INCR, which is why Redis is the recommended backend for chords. The chord body receives results in the order the group was defined, not completion order.

Step 3 — Route Failures to One Error Handler

Attach an error callback with link_error. It receives the failing task's request, the exception, and the traceback. On a chain, the errback attached to the whole chain fires for a failure in any step.

@app.task
def statement_failed(request, exc, traceback, customer_id: str) -> None:
    log.error("statement workflow failed", extra={
        "customer_id": customer_id, "task": request.task, "task_id": request.id, "error": repr(exc)})
    workflow_runs.mark_failed(customer_id, step=request.task, error=repr(exc))
    alerts.notify_if_threshold("statement_failures")

def start(customer_id: str, month: str):
    wf = statement_workflow(customer_id, month)
    return wf.apply_async(link_error=statement_failed.s(customer_id=customer_id))

Chord failures need extra care. If any group member fails permanently (retries exhausted), the chord body never runs and Celery marks the chord as failed with a ChordError. The errback fires, but the successful members' results are still sitting in the backend; if the body wrote partial state, compensate it in the errback. Error handling for single tasks is covered in Celery task retry and error handling.

One failed member cancels the join Three fetch tasks run as a chord header. The EU and US fetches succeed, but the AP fetch exhausts its retries. The chord counter never reaches three, so merge_usage is not enqueued; instead the chord is marked failed and the error callback runs. Chord header 2/3, body skipped fetch eu: success fetch us: success fetch ap: 5 retries, fails chord counter stops at 2 of 3 merge_usage: never runs statement_failed runs

Step 4 — Launch 40,000 Workflows Without Flooding the Broker

Calling start() in a loop for every customer enqueues 120,000 fetch tasks at once, pushes every other queue's jobs behind them, and creates 40,000 chord counters in Redis. Throttle the launch instead: a scheduler task starts workflows in pages and re-schedules itself.

@app.task
def launch_statements(month: str, after_id: str = "", page: int = 500) -> None:
    ids = customers.ids_after(after_id, limit=page)
    for cid in ids:
        start(cid, month)
    if ids:
        # next page in 60s: ~500 workflows/minute keeps the fetch queue shallow
        launch_statements.apply_async(args=[month, ids[-1], page], countdown=60)

Route the fetch tasks to their own queue with a dedicated worker pool so that month-end reporting cannot delay transactional work — the routing setup is described in Celery task routing with task_routes. For a single very wide fan-out (one workflow with thousands of members), use chunks to batch members into fewer tasks: fetch_usage.chunks(args_list, 100).group() turns 10,000 calls into 100 tasks of 100.

Step 5 — Configure the Result Backend for Chords

Chord reliability is result-backend reliability. Three settings matter:

# celeryconfig.py
result_backend = "redis://redis-results:6379/1"
result_expires = 3 * 24 * 3600          # must exceed your longest workflow; default 1 day
result_extended = True                   # store task name/args for debugging failed chords
redis_backend_health_check_interval = 30
task_ignore_result = False               # chord members MUST store results
worker_prefetch_multiplier = 1           # long fetches: avoid hoarding while chords wait

A task configured with ignore_result=True that is used as a chord member breaks the chord silently — the counter never increments. If you ignore results globally for performance, override it on every task that can appear in a chord header. Put the result backend on a Redis instance separate from the broker so a broker backlog cannot evict chord state, and set maxmemory-policy noeviction on it.

Step 6 — Record Each Workflow So You Can Find It Later

The AsyncResult returned by apply_async identifies only the last task of the chain. Once the Python process that launched the workflow is gone, finding "the statement workflow for customer c-17" means searching task ids in logs. Store the ids you will need when you launch, next to your own record of the run:

def start(customer_id: str, month: str) -> str:
    wf = statement_workflow(customer_id, month)
    result = wf.apply_async(link_error=statement_failed.s(customer_id=customer_id))
    run_id = workflow_runs.create(
        kind="monthly_statement",
        subject_id=customer_id,
        root_task_id=result.id,                         # id of the final task in the chain
        parent_ids=[r.id for r in iter_parents(result)], # every earlier step's id
        month=month,
    )
    return run_id

def iter_parents(result):
    node = result.parent
    while node is not None:
        yield node
        node = node.parent

With the ids stored, support tooling can show the state of every step (AsyncResult(id).state), and an operator can revoke the remaining steps of one customer's workflow without guessing. Pair the record with the progress fields described in tracking progress of multi-step jobs if users need to watch the run.

For fleet-level visibility, emit one metric per workflow completion rather than relying on per-task metrics: a histogram of end-to-end duration labelled by workflow name, and a counter of outcomes. Per-task dashboards show that fetch_usage is fast; only the workflow metric shows that statements take 40 minutes end to end because they wait behind month-end traffic between steps.

A handle on every step A workflow run record for customer c-17 stores the task ids of the three fetch tasks, the merge, the render, and the email steps. Support tooling reads the record to show the state of each step and can revoke the steps that have not run yet. workflow_runs row for c-17 monthly_statement root_task_id + parent_ids month = 2026-08 fetch x3: SUCCESS merge: SUCCESS render: STARTED email: PENDING support can inspect state and revoke pending steps

Verification

Run the workflow eagerly in a unit test to check wiring, then against a real broker to check the chord:

def test_statement_workflow_end_to_end(celery_app, celery_worker, fake_regions):
    result = statement_workflow("c-17", "2026-08").apply_async()
    assert result.get(timeout=30) is None               # email task returns None
    assert storage.exists("pdf/c-17.pdf")
    assert mailer.sent_to("c-17")

def test_chord_failure_calls_errback(celery_app, celery_worker, fake_regions):
    fake_regions["ap"].fail_always()
    start("c-18", "2026-08")
    wait_until(lambda: workflow_runs.state("c-18") == "failed", timeout=60)

In production, watch celery_task_failed_total{task="tasks.fetch_usage"} and the count of runs marked failed per hour; a spike in the second without the first means the errback is catching something other than fetch failures.

Gotchas & Edge Cases

Chords inside chords. Nested chords work but multiply result-backend traffic and make failures hard to trace. Flatten where possible: a chain of chords is easier to reason about than a chord whose body is another chord.

Large results through the backend. Returning megabytes from chord members stores them all in Redis and passes the whole list to the body. Return storage keys, not data.

Revoking a workflow. result.revoke() on a chain revokes only tasks already known; later steps are created when earlier ones finish. Store a cancelled flag in your workflow record and have each task check it before doing work.

Eager mode hides chord bugs. task_always_eager=True runs chords synchronously and never touches the result backend counter, so tests pass while production chords hang. Use the celery_worker pytest fixture for chord tests, as described in testing Celery tasks with pytest.

FAQ

Can I use a chord without a result backend? No. The chord counter and the member results live in the result backend. With result_backend unset, Celery raises NotImplementedError when you apply a chord.

Why does my chord body run twice? Usually because a member task was redelivered after it already counted (worker killed after the backend update but before the ack). Make the body idempotent, or guard it with a unique key per workflow run.

Is a chain of 20 tasks a good idea? It works, but each step adds queue wait. If steps are short and always run together, combine them into fewer tasks; keep separate tasks for steps that need different queues, retries, or resources.

Related