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.

One record, many readers Workers executing different steps of the export write phase, units done, and a heartbeat to the export run record. Web servers read the record and push updates to the browser over server-sent events. Support tooling reads the same record. Progress lives with the run, not the job query worker write worker upload worker export_runs row phase, percent, heartbeat browser (SSE) support tooling

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.

Weighting phases fixes the stall Two progress lines over the duration of an export. The unweighted line jumps to 40 percent in the first few seconds, then stays nearly flat for most of the run before jumping to 100. The weighted line rises steadily from start to finish. Reported percent over one export equal weights: stuck near 40% weighted by typical phase time start done

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.

Cancellation rides on the progress write The user presses cancel, which sets the run state to cancelled. The worker's next progress update includes a condition that the run is still running, so it updates zero rows. The worker treats that as a cancellation, deletes partial output, and acknowledges its message without enqueuing further steps. Stop within one reporting interval user: cancel state = cancelled next report: 0 rows updated worker stops, cleans up, acks No broker-side deletion needed: a queued job for a cancelled run exits at its first report.

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