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_lateredelivers a task that may have partly run. result_expireslong 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.
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.
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.
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
- Job Chaining & Workflow Orchestration — choreography, orchestration, and workflow records.
- Implementing Sagas with Compensating Jobs — undo partial work when a chord fails.
- Tracking Progress of Multi-Step Jobs — show where each workflow is.
- Tuning the Celery Result Backend — the backend that chords depend on.