skip to content

Task Canvas

Composing tasks with signatures, chain, group and chord, plus link and link_error callbacks. Interviewers ask why a chord needs a result backend and why a task must not wait on another.

on this pageshow

explore

questions

6

In Celery, how do chain, group and chord differ, and which fits resizing thumbnails in parallel and then running one step after all finish?

level: middleimportance: must knowfreq 45%

answer

  1. sequence, fan-out, fan-in
  2. each result feeds the next
  3. header plus body
  4. group piped into a task

basics

~20 s

In Celery, a chain runs tasks in sequence and passes each result on; a group runs tasks in parallel; a chord is a group plus one body task that receives the list of all header results. Parallel thumbnails, then publish, is a chord.

solid answer

~50 s

A `chain` runs signatures one after another: each step is sent only when the previous one succeeds, with its return value prepended to the next step's arguments; if a step fails, the later steps never run. A `group` sends a list of signatures at once so workers can run them in parallel, and returns a `GroupResult`. A `chord` is a group (the **header**) plus a **body** signature that runs once, after every header task has finished, with the list of their results as its first argument. For an upload that needs 160, 640 and 1280 px thumbnails and then one `publish_photo` step, the fit is `chord(make_thumbnail.s(photo_id, size) for size in sizes)(publish_photo.s(photo_id))`, or equivalently a group piped into a task, which Celery upgrades to a chord. A chord needs a result backend to know when the header is done.

code

python · 25 lines
python
from celery import Celery, chain, chord

app = Celery('photos', broker='redis://localhost:6379/0',
             backend='redis://localhost:6379/1')

@app.task
def make_thumbnail(photo_id, size):
    return f'thumbs/{photo_id}_{size}.jpg'

@app.task
def publish_photo(thumb_paths, photo_id):
    # thumb_paths is the list of all header results
    return photo_id

@app.task
def notify_owner(user_id):
    print(f'photo ready for user {user_id}')

def start_pipeline(photo_id, user_id):
    header = [make_thumbnail.s(photo_id, size) for size in (160, 640, 1280)]
    workflow = chain(
        chord(header, publish_photo.s(photo_id)),
        notify_owner.si(user_id),
    )
    return workflow.delay()

go deeper

for a junior

Recall the three shapes: chain is one after another, group is all at once, chord is all at once and then one body with every result.

for a middle

Explain how results flow: prepended into the next chain step, forwarded to every group member, collected as a list for the chord body, and why a chord needs a result backend.

for a senior

Show design judgment: when a chord's join is worth its cost, keeping header return values small, and how a failed step or header task stops the rest of the workflow.

for a principal

Discuss when a canvas workflow is the right tool at all versus a simpler single task or an external orchestrator, given the result backend and visibility it demands.

## Three ways to compose tasks Celery's **canvas** is the set of primitives that turn single task calls, called **signatures**, into workflows. Three of them cover almost every real pipeline: | Primitive | Shape | Runs | Returns | |---|---|---|---| | `chain` | A then B then C | Sequentially; each result passed on | `AsyncResult` of the last task | | `group` | A, B, C at once | In parallel, on any free workers | `GroupResult` of all tasks | | `chord` | group, then one body | Header in parallel, body once at the end | `AsyncResult` of the body | Each of them is itself a signature, so they nest: a chain can contain a chord, and a chord's body can be a chain. ## Chain: sequence with results passed forward A **chain** links signatures so that each runs after the previous one succeeds. The `|` operator builds one: `save_original.s(upload_id) | make_thumbnail.s(640)`. - Every step is a **separate task message**. Steps can run on different workers; nothing runs in one process by default. - The finished step's return value is **prepended** to the next signature's arguments, unless that signature is immutable (`.si()`). - If a step raises, the remaining steps are never sent. The failure is stored as that step's result, and calling `.get()` on the chain's final result re-raises the parent's error. - Calling a chain returns the `AsyncResult` of the **last** task; its `.parent` attribute walks back through the earlier steps. ## Group: parallel fan-out A **group** takes a list of signatures and sends them all at once. With several workers, they run in parallel. Calling a group returns a `GroupResult`, whose `.get()` returns the results in the same order as the signatures in the group, whatever order they finished in. The group itself has no step that runs "after" it; to do something once all members finish, you need a chord. When a group sits inside a chain, the previous step's result is forwarded to **every** task in the group as a partial argument. ## Chord: fan-out, then fan-in A **chord** has two parts: a **header** (a group) and a **body** (one signature). The body runs exactly once, after every header task has returned, and receives the **list of header results** as its first argument. You can write it two ways: 1. `chord(header, body)`, or `chord(header)(body)` to call it at once; 2. `group(...) | body`, because piping a group into a task makes Celery build a chord automatically. The synchronization step needs somewhere to record each header task's completion, so **a chord requires a result backend**; without one, starting it raises `NotImplementedError`. Celery's own guide calls the synchronization costly and says to use chords only where you need the join. ## The photo upload pipeline A photo-sharing app receives an upload and needs three thumbnail sizes, then one step that publishes the photo once all of them exist, then a notification to the owner: 1. **Fan out**: one `make_thumbnail` task per size, in parallel (the chord header). 2. **Fan in**: `publish_photo` receives the list of thumbnail paths and marks the photo live (the chord body). 3. **Continue**: `notify_owner` runs after publishing (a chain after the chord). A chain alone would resize the sizes one at a time. A group alone would resize them in parallel but give you no place to publish when they are all done. The chord is the primitive that expresses "all of these, then that". ## What to say about the trade-offs - A chain's latency is the **sum** of its steps; a chord's header latency is its **slowest** member, plus the join. - A chord body sees **all** header results at once, so a header with thousands of large return values means a large body message and backend load. Return paths, not image bytes. - A failure in one header task means the body never runs; that is a separate topic (chord error handling) with its own knobs. - Groups and chords give **no ordering** of execution; only the result list is ordered. Interviewers often follow up by asking what the caller gets back. For the pipeline above, calling the outer chain returns the `AsyncResult` of `notify_owner`, the last step, and its `.parent` attribute walks back toward the earlier steps. Calling a chord on its own returns the body's `AsyncResult`, not the header's. In one sentence: chain for "then", group for "at the same time", chord for "at the same time, then once".

  • What does Celery do with group(...) | publish_photo.s(photo_id)?
    Piping a group into a signature makes Celery build a `chord`: the group becomes the header and the signature the body. So it inherits chord requirements, above all a configured result backend. Celery's own error message notes that a group chained with a task is upgraded to a chord because the pattern needs synchronization.
  • In a Celery chain, what happens to later steps if the middle step raises?
    They are never sent. The failing task's result holds the exception, and the next signature is not applied because it is only sent on success. Calling `.get()` on the chain's final `AsyncResult` re-raises the parent's error by default. To react to the failure, attach an errback with `link_error` or `on_error()`; on a chain it is linked to every step.
  • What happens if a chord's header is empty, for example a photo with no thumbnail sizes requested?
    Nothing in the header would ever complete, so Celery does not wait for it: with no header tasks, `chord.run` sends the body at once with an empty list, `body.delay([])`. The body must therefore handle an empty list of results rather than assume at least one.

saying these in an interview costs you the question

  • A chain runs all its tasks inside one worker process as one message.
  • A chord body runs each time a header task finishes.
  • Piping a group into a task runs that task once per group member.
  • A group guarantees its tasks execute in the order they are listed.
  • When a chain step fails, Celery skips it and continues with the next step.
open as a page

Why does calling AsyncResult.get() inside a Celery task raise RuntimeError, and how should a thumbnail workflow wait for another task instead?

level: middleimportance: must knowfreq 35%

basics

~20 s

Celery's AsyncResult.get() defaults to disable_sync_subtasks=True, so inside a worker task it raises RuntimeError: a task blocking on another can exhaust the pool and deadlock. Express the dependency with a chain, chord or link callback instead.

open as a page

In Celery, what is a task signature, and how do .s() and .si() differ when the signature runs inside a chain?

level: juniorimportance: should knowfreq 38%

basics

~20 s

A Celery signature is one task call packed as data: task name, args, kwargs and execution options. In a chain, a .s() signature gets the previous step's result prepended to its arguments; an immutable .si() signature ignores it.

open as a page

Why does a Celery chord need a result backend, and what happens to the chord body when one thumbnail task in its header fails?

level: seniorimportance: should knowfreq 28%

basics

~20 s

A Celery chord needs a result backend because something must record each header task's completion and result before the body can start. If one header task fails, the others still run, but the body never runs and is marked FAILURE with a ChordError.

open as a page

In Celery, how do group, chunks() and starmap() differ when resizing 10,000 archived photos, in task messages sent and in parallelism?

level: middleimportance: nice to knowfreq 16%

basics

~20 s

In Celery, a group sends one message per call, so 10,000 photos means 10,000 parallel tasks. starmap() sends one message and runs every call in sequence inside one task. chunks(it, n) sends a group of starmap tasks, each running n calls in sequence.

open as a page