Reporting Background Job Progress and Status¶
A request that starts a five-minute export cannot wait for it, so the job runs in the background and the user gets a job id. The next questions are always the same: is it running, how far has it got, is it finished, can I cancel it? A job system answers the first and third out of the box; progress and cancellation need a little help. Tested with arq 0.28 and Redis 7.4: polling job.status() showed queued, then in_progress, then complete; a progress hash the job updated every tenth row showed 90 of 100 on the last poll before completion; and job.abort() on a running 1,000-row export returned True after 0.5 s. This guide wires those pieces into an API a front end can use.
Prerequisites¶
- Python 3.11+,
pip install arqand Redis; a web framework for the endpoints (FastAPI here). - arq workers, from running arq workers with Redis.
- Streaming to the browser, from streaming server-sent events from asyncio.
1. Read status from the job system¶
arq stores each job's state in Redis, and a Job handle reads it:
from arq.jobs import Job, JobStatus
async def job_state(arq_redis, job_id: str) -> dict:
job = Job(job_id, arq_redis)
status = await job.status() # deferred, queued, in_progress, complete, not_found
out = {"id": job_id, "status": status.value}
if status == JobStatus.complete:
info = await job.result_info()
out["success"] = info.success
out["result"] = info.result if info.success else repr(info.result)
return out
Observed over a 2-second export: queued before the worker picked it up, in_progress while running, complete afterwards with the result attached. not_found means the id never existed or its result has expired — keep_result on the worker controls how long completed results stay, so set it longer than any client will plausibly poll.
Verify: query the state of a job before, during and after it runs; the three statuses appear in order.
2. Publish progress from inside the job¶
The job system does not know how far a job has got; the job does. Write progress to a Redis hash keyed by job id, with a TTL so abandoned keys clean themselves up:
async def export_orders(ctx, account_id: int) -> dict:
r = ctx["redis"]
key = f"progress:{ctx['job_id']}"
total = await count_orders(account_id)
done = 0
async for batch in fetch_orders(account_id, batch_size=500):
await write_rows(batch)
done += len(batch)
await r.hset(key, mapping={"done": done, "total": total, "stage": "writing"})
await r.expire(key, 3600)
await r.hset(key, mapping={"done": done, "total": total, "stage": "finished"})
return {"rows": done}
Update per batch, not per row. Each update is a Redis round trip of well under a millisecond, but a job processing 100,000 rows with a write per row would spend measurable time on progress alone; one write per few hundred rows or per second is plenty for a progress bar. Include a stage field for multi-phase jobs — "fetching", "writing", "uploading" — because a percentage that stalls at 100% during upload confuses users more than a named stage.
Verify: poll the hash during a run; done increases monotonically and reaches total.
3. Expose status and progress through one endpoint¶
Clients should not need to know about Redis keys. Combine both into one response:
from fastapi import FastAPI, HTTPException
app = FastAPI()
@app.post("/exports")
async def start_export(account_id: int):
job = await app.state.arq.enqueue_job("export_orders", account_id)
return {"job_id": job.job_id}
@app.get("/exports/{job_id}")
async def export_status(job_id: str):
state = await job_state(app.state.arq, job_id)
if state["status"] == "not_found":
raise HTTPException(404, "unknown or expired job")
progress = await app.state.arq.hgetall(f"progress:{job_id}")
if progress:
done, total = int(progress[b"done"]), int(progress[b"total"]) or 1
state["progress"] = {"done": done, "total": total, "pct": round(100 * done / total),
"stage": progress.get(b"stage", b"").decode()}
return state
Authorise the read: a job id in a URL is not a secret, so check that the requesting user owns the job — store the owner in the progress hash or a separate key at enqueue time.
Verify: a client polling once a second sees status and percentage move together, and another user's request for the same id gets a 403 or 404.
4. Push progress instead of polling, when it matters¶
Polling once a second from many clients is fine for most products. For live dashboards, push updates with server-sent events, publishing from the job and relaying in the API:
from fastapi.responses import StreamingResponse
async def export_orders(ctx, account_id: int) -> dict:
...
await r.publish(f"progress:{ctx['job_id']}", f"{done}/{total}") # alongside the HSET
@app.get("/exports/{job_id}/events")
async def export_events(job_id: str):
async def stream():
async with app.state.arq.pubsub() as ps:
await ps.subscribe(f"progress:{job_id}")
async for msg in ps.listen():
if msg["type"] == "message":
yield f"data: {msg['data'].decode()}\n\n"
return StreamingResponse(stream(), media_type="text/event-stream")
Keep the hash as the source of truth and the channel as a notification: a client that connects mid-job reads the hash for the current value, then listens for changes. Pub/sub does not replay missed messages. Each open SSE connection holds a Redis subscription, so cap them per user.
Verify: a browser receives progress events in real time, and reconnecting mid-job shows the current value immediately.
5. Let users cancel a running job¶
arq can abort a running job when the worker is started with allow_abort_jobs=True; the job is cancelled at its next await:
class WorkerSettings:
functions = [export_orders]
allow_abort_jobs = True
@app.delete("/exports/{job_id}")
async def cancel_export(job_id: str):
job = Job(job_id, app.state.arq)
aborted = await job.abort(timeout=5)
return {"aborted": aborted}
Measured: aborting a running export returned True after 0.5 s, the time for the worker to notice the abort request on its next poll. The job sees asyncio.CancelledError, so the cleanup rules for cancellation apply — release locks, delete partial files, update the progress hash to "cancelled" in a finally block, as described in preventing CancelledError leaks in cleanup.
Verify: cancel mid-run; the job stops within a second, partial output is cleaned up, and status reports the job as finished without a successful result.
Verification¶
Progress reporting is complete when:
- Status comes from the job system and distinguishes queued, running, complete and expired.
- Progress is written per batch to a hash with a TTL, with a stage name for multi-phase jobs.
- One authorised endpoint returns both to clients.
- Cancellation works, and cancelled jobs clean up and report a final state.
Diagnostic Hook: export the age of the oldest in_progress job per type, computed from the job's start time, and the number of progress keys whose done has not changed for several minutes. A running job with stalled progress is stuck on something — typically an external call without a timeout — and is worth alerting on before the user gives up.
Pitfalls & edge cases¶
keep_resultshorter than client polling. Completed jobs turn intonot_found, which looks like an error.- Progress writes per row. Turns a progress bar into a significant share of the job's runtime.
- Unauthenticated status endpoints. Job ids leak through logs and URLs; check ownership.
- Pub/sub as the only source. Clients that connect mid-job see nothing until the next message.
Frequently Asked Questions¶
How do I check the status of an arq job?
Create arq.jobs.Job(job_id, redis) and await job.status(), which returns queued, deferred, in_progress, complete or not_found. For completed jobs, job.result_info() returns the result and whether it succeeded.
How do I report progress from a background job?
Have the job write done, total and a stage name to a Redis hash keyed by its job id, once per batch, with a TTL. An API endpoint reads that hash together with the job status.
Should clients poll or stream job progress?
Polling once a second is simplest and enough for most interfaces. Use server-sent events with Redis pub/sub when updates must be live, keeping the hash as the source of truth for clients that connect mid-job.
How do I cancel a running arq job?
Start the worker with allow_abort_jobs=True and call await job.abort(). The job is cancelled at its next await, so its cleanup must handle CancelledError.
Related¶
- Background Jobs & Task Queues — up to the topic overview.
- Checkpointing progress in long-running async jobs — progress that also lets a job resume.
- Concurrent Execution & Worker Patterns — the section overview.