Tracking Progress of Multi-Step Jobs
Users tolerate a slow background job far better than an opaque one, and this guide shows how to report honest progress for multi-step work as part of Job Chaining & Workflow Orchestration in Queue Fundamentals & Architecture. It builds a progress record that survives worker restarts, computes a percentage that does not jump backwards, streams updates to the browser, and flags runs that have silently stopped moving.
Problem Statement
A data-export feature runs as a workflow: query the account's records in pages, write them to CSV chunks, zip the chunks, upload the archive, and email a link. Exports take between 20 seconds and 40 minutes. The UI shows a spinner. Support tickets ask "is my export still running?", and when a worker crashes the spinner spins forever. You want a progress bar that moves smoothly and truthfully, survives restarts, reaches a final "done" or "failed" state, and lets support see where any export is.
Prerequisites
- A workflow that already runs as separate steps or a single long job with identifiable phases (see Job Chaining & Workflow Orchestration).
- A database table you can write to from workers (Postgres here), and optionally Redis for pub/sub fan-out to web servers.
- A web framework that can hold a streaming response open (server-sent events) or a client that polls.
Step 1 — Create a Progress Record Owned by the Run
Progress belongs to the workflow run, not to any single job. Store it in a table keyed by run id so any worker and any web server can read it, and so it survives broker restarts and result-backend expiry.
CREATE TABLE export_runs (
id uuid PRIMARY KEY,
account_id text NOT NULL,
state text NOT NULL DEFAULT 'queued', -- queued|running|done|failed|cancelled
phase text NOT NULL DEFAULT 'queued', -- query|write|zip|upload|notify
done_units bigint NOT NULL DEFAULT 0, -- records exported so far
total_units bigint, -- NULL until known
percent numeric(5,2) NOT NULL DEFAULT 0,
message text,
heartbeat_at timestamptz NOT NULL DEFAULT now(),
created_at timestamptz NOT NULL DEFAULT now(),
finished_at timestamptz
);
CREATE INDEX export_runs_running ON export_runs (heartbeat_at) WHERE state = 'running';
The heartbeat_at column is what distinguishes "slow" from "dead" in Step 5. total_units starts as NULL because the record count is not known until the query phase has counted — the UI should show an indeterminate state until then rather than a fake percentage.
Step 2 — Weight the Phases So the Bar Moves Honestly
A naive progress bar gives each of five phases 20%. In practice the query-and-write phases take 90% of the time, so the bar races to 40%, crawls, then jumps to 100%. Weight phases by their typical share of wall-clock time, measured from past runs, and report progress within the current phase from units done.
# progress.py
PHASES = [ # (name, weight) — weights from p50 durations of past runs
("query", 0.10),
("write", 0.70),
("zip", 0.10),
("upload", 0.08),
("notify", 0.02),
]
_START = {}
acc = 0.0
for name, w in PHASES:
_START[name] = acc
acc += w
def overall_percent(phase: str, done: int, total: int | None) -> float:
weight = dict(PHASES)[phase]
within = (done / total) if total else 0.0
return round(100 * (_START[phase] + weight * min(within, 1.0)), 2)
Recompute the weights occasionally from completed runs (avg(phase_seconds) / avg(total_seconds) per phase). They do not need to be exact; they need to prevent the bar from standing still for most of the run.
Step 3 — Write Updates Without Hammering the Database
A write worker exporting 2 million rows must not issue 2 million UPDATEs. Throttle updates by time and by change, and never let progress go backwards — a retried chunk can briefly report a lower done_units than the one before it.
import time
class ProgressReporter:
def __init__(self, run_id: str, min_interval: float = 1.0, min_delta: float = 0.5):
self.run_id, self.min_interval, self.min_delta = run_id, min_interval, min_delta
self._last_write, self._last_pct = 0.0, -1.0
def report(self, phase: str, done: int, total: int | None, message: str = "") -> None:
pct = overall_percent(phase, done, total)
now = time.monotonic()
if now - self._last_write < self.min_interval and pct - self._last_pct < self.min_delta:
return
db.execute(
"""UPDATE export_runs
SET phase=:phase, done_units=:done, total_units=COALESCE(:total, total_units),
percent=GREATEST(percent, :pct), message=:msg,
heartbeat_at=now(), state='running'
WHERE id=:id AND state IN ('queued','running')""",
id=self.run_id, phase=phase, done=done, total=total, pct=pct, msg=message)
redis.publish(f"export:{self.run_id}", f"{pct}:{phase}") # wake SSE listeners
self._last_write, self._last_pct = now, pct
GREATEST(percent, :pct) makes the bar monotonic even when a retried chunk reports stale numbers, and the state IN (...) guard stops a late update from resurrecting a run that was already cancelled or failed. One write per second per active export is negligible for Postgres even with thousands running.
Step 4 — Stream Progress to the Browser
Polling every two seconds works and is the right first version. Server-sent events feel better and cost less at scale: the web server subscribes to the run's Redis channel and forwards each update, falling back to the database row on connect so a late subscriber sees the current state immediately.
# views.py — FastAPI server-sent events
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import json
app = FastAPI()
@app.get("/exports/{run_id}/events")
async def export_events(run_id: str):
async def stream():
row = await db.fetch_one("SELECT state, phase, percent, message FROM export_runs WHERE id=$1", run_id)
yield f"data: {json.dumps(dict(row))}\n\n" # current state first
if row["state"] in ("done", "failed", "cancelled"):
return
pubsub = aioredis.pubsub()
await pubsub.subscribe(f"export:{run_id}")
async for msg in pubsub.listen():
if msg["type"] != "message":
continue
row = await db.fetch_one("SELECT state, phase, percent, message FROM export_runs WHERE id=$1", run_id)
yield f"data: {json.dumps(dict(row))}\n\n"
if row["state"] in ("done", "failed", "cancelled"):
break
return StreamingResponse(stream(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})
The pub/sub message is only a wake-up; the data always comes from the database row, so a missed pub/sub message (Redis pub/sub is fire-and-forget) costs one update's worth of staleness rather than a wrong display.
Step 5 — Detect Runs That Stopped Moving
A crashed worker stops sending heartbeats. A sweeper marks runs whose heartbeat is older than the longest reasonable gap between updates, and either re-enqueues the current phase or fails the run with a clear message.
STALE_AFTER = "5 minutes"
@app_celery.task
def sweep_stale_exports() -> None:
stale = db.fetch_all(
f"""UPDATE export_runs SET state='failed',
message='Export stopped responding; please retry', finished_at=now()
WHERE state='running' AND heartbeat_at < now() - interval '{STALE_AFTER}'
RETURNING id, account_id, phase""")
for run in stale:
metrics.exports_stalled.labels(phase=run.phase).inc()
redis.publish(f"export:{run.id}", "failed") # close open SSE streams
Choose between failing and resuming per phase: resuming the write phase from the last completed chunk is valuable for a 40-minute export; for a 20-second one, failing and letting the user retry is simpler. The phase label on the metric shows which phase most often stalls — usually the one with the largest memory footprint, which points at an OOM kill rather than a bug.
Step 6 — Let Users Cancel from the Progress View
Once users can see a long export crawling, they will want to stop it. Cancellation from the UI is a state change on the same record, and workers honour it at their next progress report — which is conveniently already happening every second.
# views.py — cancel endpoint
@app.post("/exports/{run_id}/cancel")
async def cancel_export(run_id: str, user=Depends(current_user)):
changed = await db.execute(
"""UPDATE export_runs SET state='cancelled', message='Cancelled by user',
finished_at=now()
WHERE id=$1 AND account_id=$2 AND state IN ('queued','running')""",
run_id, user.account_id)
if changed:
await aioredis.publish(f"export:{run_id}", "cancelled") # close open streams
return {"cancelled": bool(changed)}
# worker side: ProgressReporter.report() already filters on state IN ('queued','running').
# Make it tell the caller when the update did not apply, and stop.
class ExportCancelled(Exception):
pass
def report_or_stop(reporter, *args):
if reporter.report(*args) is False: # UPDATE matched 0 rows: cancelled or failed
raise ExportCancelled()
The worker catches ExportCancelled, deletes any partial archive it wrote, acknowledges its message, and enqueues nothing further. Because the check rides on the existing progress write, cancellation costs no extra queries and takes effect within one reporting interval. The account_id filter on the endpoint matters: a run id in a URL must never let one customer cancel another's export.
Queued runs that have not started yet are handled by the same filter — the worker that eventually picks up the job finds the run already cancelled at its first report and exits immediately. There is no need to find and delete the job message from the broker, which many brokers cannot do efficiently anyway.
Verification
def test_progress_is_monotonic_and_terminal(run_export, kill_worker_during):
run_id = run_export(account_id="a1", rows=50_000)
samples = poll_percent(run_id, every=0.2)
assert samples == sorted(samples) # never goes backwards
assert db.scalar("SELECT state FROM export_runs WHERE id=:i", i=run_id) == "done"
def test_crashed_export_becomes_failed(run_export, kill_worker_during, freeze_time):
run_id = run_export(account_id="a2", rows=500_000)
kill_worker_during("write")
freeze_time.advance(minutes=6)
sweep_stale_exports()
assert db.scalar("SELECT state FROM export_runs WHERE id=:i", i=run_id) == "failed"
In production, count(state='running' AND heartbeat_at < now() - interval '2 minutes') should be near zero; a sustained non-zero count means workers are wedged rather than crashed.
Gotchas & Edge Cases
Progress from the result backend. Celery's update_state(meta=...) stores progress in the result backend, which expires results and may be evicted. It is fine for a toy; for anything users depend on, write to your own table.
Counting before querying. SELECT count(*) on a large table to get total_units can take longer than the export itself. Use an estimate (pg_class.reltuples or a cached count) and correct it as you go, or show indeterminate progress for the first phase.
Fan-out phases. When a phase runs as parallel chunks, sum done_units across chunks with an atomic increment (UPDATE ... SET done_units = done_units + :n), not by each chunk writing its own count.
Time remaining estimates. Estimated time remaining from the recent rate is far more accurate than from overall average; use the last one to two minutes of progress and clamp wild swings.
FAQ
Polling or server-sent events? Start with polling every 2–5 seconds; it is simple and cache-friendly. Move to SSE when many users watch progress at once or when the smoothness matters. WebSockets are unnecessary for one-way progress.
Should progress updates go through the job queue? No. Progress is state, not work. Writing it to a table (plus an optional pub/sub nudge) is cheaper and never delays real jobs.
How do I show progress for a Celery chord? Count completed members in the run record with an atomic increment in each member task, and derive the percentage from completed over total members. The chord itself does not expose progress.
Related
- Job Chaining & Workflow Orchestration — the workflow record this progress extends.
- Celery Chains, Groups, and Chords — framework primitives whose progress you may need to report.
- Alerting on Stuck and Stalled Jobs — the operator-side view of stalled work.
- Configuring Visibility Timeouts for Long-Running Workers — heartbeats at the broker level.