Job Chaining & Workflow Orchestration
Most real background work is not one job but a sequence of them — transcode, then thumbnail, then notify — and this guide covers how to compose jobs into workflows as part of Queue Fundamentals & Architecture. The moment a second job depends on the first, you inherit questions that a single job never raises: where does the intermediate result live, what happens to step three when step two fails, how does anyone see that the workflow is 60% done, and who cleans up after a half-finished run.
Queues handle short chains well and long, branching, human-in-the-loop processes badly. Knowing where that line sits — and building chains that fail cleanly on the right side of it — is what separates a pipeline that recovers on its own from one that leaves orphaned half-results in the database every time a worker restarts.
The Scenario: A Pipeline That Leaves Debris
A video platform processes uploads in four steps: probe the file, transcode three renditions in parallel, generate a thumbnail sprite, then mark the video as ready and email the uploader. The first implementation had each job enqueue the next at the end of its body. It worked in testing. In production, three kinds of debris appeared:
- Orphans. A transcode worker was OOM-killed after writing two of three renditions. The job retried and succeeded, but the "all renditions done" check in the last transcode ran twice — once per surviving attempt — and the thumbnail job was enqueued twice.
- Dead ends. A thumbnail job failed permanently on a corrupt frame. Nothing downstream ever ran, the video stayed in
processingforever, and nobody was alerted because no individual job was "stuck" — the workflow just stopped. - Invisible progress. Support could not answer "where is my video?" because the state was spread across four queues and the database, with no single record of which step a given upload had reached.
Each problem has a standard fix: idempotent fan-in, an explicit workflow record with a terminal state, and compensation for steps that must be undone. The rest of this guide builds those pieces.
Architectural Overview: Choreography vs Orchestration
There are two ways to connect jobs, and most systems use both in different places.
Choreography means each job knows its successor and enqueues it. It needs no central component and is the natural first implementation. Its weakness is that the workflow exists only as a chain of side effects: no record says "this upload is at step 3 of 5", and a failure in the middle leaves no one responsible for the rest.
Orchestration means a coordinator — a workflow record plus a small state machine, a framework primitive like Celery's canvas or BullMQ flows, or a dedicated engine like Temporal — owns the sequence. Jobs report completion to it, and it decides what runs next. The workflow is visible, has a terminal state, and can be retried or cancelled as a unit.
# workflow-definition.yaml — the explicit shape the orchestrator enforces
workflow: video-upload
version: 3
steps:
- id: probe
queue: media-light
timeout: 60s
- id: transcode
queue: media-heavy
fan_out: [1080p, 720p, 480p] # one job per rendition
timeout: 30m
retries: 3
- id: thumbnail
after: transcode # fan-in: waits for all three
queue: media-light
- id: notify
after: thumbnail
queue: default
on_failure:
compensate: [delete_renditions, release_storage_quota]
terminal_state: failed # every run ends in done | failed | cancelled
deadline: 2h # the whole workflow, not per step
The last two fields are where choreography usually falls short. An orchestrated workflow has a deadline for the whole run and a guaranteed terminal state, so "stuck in processing forever" becomes a detectable, alertable condition.
Implementation 1: Framework Primitives for Chains and Fan-In
Every major job framework ships workflow primitives, and for pipelines of a handful of steps they are the right starting point. Celery's canvas composes chain, group, and chord:
# pipeline.py — Celery canvas
from celery import chain, chord, group
from tasks import probe, transcode, thumbnail, notify, mark_failed
def start_pipeline(video_id: str):
workflow = chain(
probe.si(video_id),
chord(
group(transcode.si(video_id, r) for r in ("1080p", "720p", "480p")),
thumbnail.si(video_id), # chord body: runs once after all three
),
notify.si(video_id),
)
return workflow.apply_async(link_error=mark_failed.si(video_id))
The chord's fan-in counter lives in the result backend, which means a chord requires a result backend and inherits its reliability — the details and failure behaviours are covered in Celery chains, groups, and chords. BullMQ has an equivalent with FlowProducer, where a parent job waits for its children, described in BullMQ flows for parent-child jobs; Sidekiq Pro's batches provide callbacks when a set of jobs completes, covered in Sidekiq batch jobs and workflows.
Framework primitives share one limitation: the workflow's state is stored in framework internals (result backend keys, Redis hashes) rather than in your database. That makes it hard to query "all uploads stuck at thumbnail for more than an hour", and it means an operator who flushes Redis also deletes every in-flight workflow.
Implementation 2: An Explicit Workflow Record
When visibility and recovery matter more than brevity, keep the workflow state in your own database and let each job report back. A single table plus an atomic fan-in counter covers most pipelines:
-- workflow_runs.sql
CREATE TABLE workflow_runs (
id uuid PRIMARY KEY,
workflow text NOT NULL,
subject_id text NOT NULL, -- e.g. video id
state text NOT NULL DEFAULT 'running', -- running | done | failed | cancelled
current_step text NOT NULL,
pending int NOT NULL DEFAULT 0, -- outstanding fan-out jobs
started_at timestamptz NOT NULL DEFAULT now(),
deadline_at timestamptz NOT NULL,
updated_at timestamptz NOT NULL DEFAULT now()
);
CREATE TABLE workflow_steps (
run_id uuid REFERENCES workflow_runs(id),
step text,
shard text, -- rendition for fan-out steps
state text NOT NULL DEFAULT 'pending',
PRIMARY KEY (run_id, step, shard) -- makes completion idempotent
);
# fan_in.py — exactly-once advance, however many times a transcode retries
def complete_transcode(run_id: str, rendition: str) -> None:
with db.transaction() as tx:
changed = tx.execute(
"UPDATE workflow_steps SET state='done' "
"WHERE run_id=:r AND step='transcode' AND shard=:s AND state<>'done'",
r=run_id, s=rendition,
).rowcount
if changed == 0:
return # duplicate completion: ignore
remaining = tx.scalar(
"UPDATE workflow_runs SET pending = pending - 1, updated_at = now() "
"WHERE id=:r RETURNING pending", r=run_id,
)
if remaining == 0:
tx.execute("UPDATE workflow_runs SET current_step='thumbnail' WHERE id=:r", r=run_id)
tx.enqueue_after_commit("thumbnail", run_id=run_id) # outbox, not a direct publish
The state<>'done' guard is what fixes the duplicate-fan-in bug from the scenario: a retried transcode that completes a second time changes no row, so the counter cannot be decremented twice. Enqueuing through an outbox after commit means the thumbnail job is published if and only if the transaction committed — the same guarantee described in the transactional outbox pattern. Exposing this record through an API also answers the support question; tracking progress of multi-step jobs builds that out.
Passing Data Between Steps
Every step of a workflow produces something the next step needs, and how that output travels is one of the first design decisions that bites. There are three options, and only one of them scales.
Through the broker. The step returns its output and the framework passes it to the next task as an argument — Celery chains do this by default. It is convenient for small values such as an id or a status string. It is a trap for anything large: a 15 MB transcoding manifest passed through three steps is serialized, stored, and deserialized three times, sits in broker memory while queued, and may exceed the broker's message limit outright (256 KB on SQS, 512 MB by default on Redis but with painful latency long before that).
Through the result backend. The step stores its output in the result store, and the next step fetches it. This is what a chord body does with its members' results. It keeps the broker light, but result backends are caches: they expire entries (result_expires), they can be evicted under memory pressure, and they are rarely backed up. A workflow that resumes after a long wait can find its inputs gone.
Through durable storage, by reference. The step writes its output to object storage or a database table and passes only the key. Every step becomes restartable from its inputs, intermediate artifacts can be inspected during an incident, and the broker only ever carries small messages.
# claim-check between steps: each step reads by key and writes by key
@app.task(acks_late=True)
def transcode(video_id: str, rendition: str) -> str:
src = storage.get(f"uploads/{video_id}/source.mp4")
out_key = f"renditions/{video_id}/{rendition}.mp4"
if not storage.exists(out_key): # idempotent: skip finished work
storage.put(out_key, ffmpeg_transcode(src, rendition))
return out_key # a short string, not the bytes
@app.task(acks_late=True)
def thumbnail(rendition_keys: list[str], video_id: str) -> str:
best = max(rendition_keys, key=resolution_of) # reads only what it needs
return storage.put(f"thumbs/{video_id}/sprite.jpg", make_sprite(storage.get(best)))
Write outputs under deterministic keys derived from the run and step, as above. A retried step then overwrites or skips the same object rather than creating a second copy, and cleanup after a failed run is a prefix delete. Set a lifecycle rule on the intermediate prefix so abandoned artifacts expire after a week.
Cancellation and Deadlines Across Steps
Users cancel exports, operators kill runaway backfills, and upstream systems withdraw requests. A queue has no notion of cancelling a workflow: revoking one task does nothing about the tasks that will be enqueued when it would have finished, and tasks already running ignore revocation unless the worker is told to terminate them.
The dependable mechanism is a cancellation flag on the workflow record, checked at step boundaries and at safe points inside long steps:
class Cancelled(Exception):
pass
def check_cancelled(run_id: str) -> None:
if db.scalar("SELECT state FROM workflow_runs WHERE id=:r", r=run_id) == "cancelled":
raise Cancelled(run_id)
@app.task(acks_late=True)
def transcode(run_id: str, video_id: str, rendition: str) -> str:
check_cancelled(run_id) # before starting expensive work
for segment in segments(video_id):
check_cancelled(run_id) # between segments of a long step
encode_segment(segment, rendition)
return finalize(video_id, rendition)
A step that raises Cancelled should acknowledge its message and enqueue nothing further; the run's compensations (delete partial renditions, release quota) then run exactly as they would for a failure. Deadlines work the same way from the other direction: the sweeper sets the flag when deadline_at passes, and every step sees it at its next check. Checking the flag costs a primary-key lookup — cache it for a few seconds inside tight loops if it shows up in profiles.
Trade-off Analysis
| Approach | Visibility | Failure handling | Long waits (hours/days) | Operational cost |
|---|---|---|---|---|
| Choreography (job enqueues next) | None beyond logs | Each job on its own; dead ends are silent | Poor | Lowest |
| Framework primitives (chord, flows, batches) | Framework tooling only | Error callbacks per workflow | Poor: state in broker/result store | Low |
| Workflow record in your DB | Full, queryable with SQL | Explicit terminal states, sweepers | Reasonable with a scheduler | Medium: you own the state machine |
| Durable workflow engine (Temporal, Step Functions) | Full history per run | Built-in retries, timeouts, compensation | Excellent: timers are first-class | Highest: new infrastructure |
The deciding question is usually how long a workflow waits. Anything that sleeps for days, waits on a human approval, or needs per-step timers and cancellation is re-implementing a workflow engine on top of a queue. When to move from job queues to Temporal walks through the signals.
Failure Modes & Recovery
The double fan-in. Retries of parallel steps each believe they were last and trigger the next step twice. Remediation: count completions with an idempotent, per-shard update as above, never with "check if all outputs exist" logic that can race.
The silent dead end. A step fails permanently and nothing downstream runs. Remediation: every run has a deadline and a sweeper that marks overdue runs failed and alerts. The sweeper is a scheduled job that runs SELECT id FROM workflow_runs WHERE state='running' AND deadline_at < now().
Partial side effects. Step 3 of 4 fails after steps 1 and 2 charged a card and reserved stock. Retrying the workflow from scratch charges twice; abandoning it leaves money and stock held. Remediation: define a compensating action for every step with an external side effect and run them in reverse on failure — the saga pattern, detailed in implementing sagas with compensating jobs.
Version skew mid-run. You deploy a new workflow definition while 3,000 runs are in flight. A run that started under version 2 resumes under version 3 and hits a step that no longer exists. Remediation: store the definition version on the run and keep old step handlers until every run on that version has finished.
Performance Tuning
- Keep payloads out of the chain. Passing a 20 MB intermediate result from one step to the next through the broker multiplies memory and serialization cost at every hop. Write intermediates to object storage and pass keys — the claim-check pattern.
- Separate queues per step weight. Put light steps (probe, notify) and heavy steps (transcode) on different queues with different worker pools, so a burst of transcodes does not delay the notify step of workflows that are otherwise finished.
- Cap fan-out width. A
groupof 10,000 jobs enqueues 10,000 messages at once and can starve every other workflow. Chunk large fan-outs (Celery'schunks, or batches of a few hundred) and let the next chunk start when the previous one drains. - Index the sweeper. The deadline query runs every minute; a partial index
ON workflow_runs (deadline_at) WHERE state = 'running'keeps it cheap as the table grows. - Measure end-to-end, not per step. A workflow's user-visible latency is the sum of queue wait and run time across every step. Record
started_atand completion time on the run and graph the percentile — per-step metrics alone hide queue waits between steps.
# p95 end-to-end workflow duration, from a histogram emitted when runs finish
histogram_quantile(0.95,
sum by (le, workflow) (rate(workflow_run_duration_seconds_bucket{state="done"}[30m])))
# Runs past their deadline (should be zero)
workflow_runs_overdue{workflow="video-upload"} > 0
FAQ
Should each job enqueue the next one? For two or three steps with no fan-in and no compensation, it is fine. As soon as you need parallel steps that join, a record of progress, or cleanup on failure, move the sequencing into a coordinator — a framework primitive or a workflow table.
Where should intermediate results live? In durable storage you control — object storage for large artifacts, your database for small ones — referenced by id in the job payload. The broker and the result backend are transport, not storage; both may evict or expire data before the workflow ends.
How do I retry a whole workflow? Only if every step is idempotent or compensated. The safer approach is to resume from the failed step using the workflow record, skipping steps already marked done.
How do I test a multi-step workflow? Test each step in isolation for its own logic, then run the whole workflow against a real broker and worker in CI with failures injected at every step boundary: a permanent failure in each step, a worker killed after each external call, and a duplicate delivery of each completion message. Assert on the final state of the workflow record and on the side effects left behind, not on the order of log lines. In-process eager modes skip the broker and hide exactly the fan-in and redelivery bugs these tests exist to catch.
Is Celery canvas production-ready for large fan-outs? It works, but chords with thousands of members put heavy load on the result backend and fail badly if a member's result expires. For wide fan-outs, prefer a counter in your own database or chunked groups.
Related
- Celery Chains, Groups, and Chords — framework primitives for sequential and parallel steps.
- Implementing Sagas with Compensating Jobs — undoing partial side effects cleanly.
- Tracking Progress of Multi-Step Jobs — expose workflow state to users and support.
- When to Move from Job Queues to Temporal — the signals that a workflow engine is worth it.
- Message Ordering Guarantees — sequencing within a stream rather than across steps.