Streaming Background Job Progress with SSE Permalink to this section
Part of Progress Streaming for Long-Running Jobs, under Real-Time Application Patterns.
Queue libraries already track progress. BullMQ has job.updateProgress(), Celery has update_state(). What they do not do is push that progress to a browser. This guide connects the two ends — a real worker in either ecosystem and a Server-Sent Events endpoint — without polling the queue, and with a stream that ends by itself when the job does.
Symptom & Developer Intent Permalink to this section
Teams usually arrive here from one of these starting points:
- The frontend polls
GET /jobs/:idevery two seconds, and the job-status endpoint is now the busiest route in the API. - Progress is streamed, but the SSE handler polls the queue’s job record internally in a loop, so every open stream costs a Redis read per second.
- Progress appears in bursts and the bar jumps, because the worker only reports at stage boundaries.
- The stream stays open after the job fails, and the page shows a spinner forever.
The intent is a worker that reports progress at a steady cadence, an SSE endpoint that reacts to those reports rather than polling, and a clean end in every terminal case — with the same design in Node.js and Python.
Root Cause Analysis Permalink to this section
The queue’s job record is the right place for state, but it is a pull interface. Both BullMQ and Celery store progress in Redis (or another result backend) and expect you to read it. An SSE handler that reads it in a loop turns every viewer into a poller, just moved from the browser to the server.
BullMQ’s QueueEvents class consumes a Redis stream that the queue writes on every progress update, completion and failure. That is exactly the push feed an SSE endpoint needs. Celery’s update_state writes to the result backend and notifies nobody, so Celery workers need a small helper that publishes to a Redis channel alongside the state write.
Step-by-Step Resolution Permalink to this section
Step 1 — Report progress from a BullMQ worker Permalink to this section
// worker.js — BullMQ
import { Worker } from 'bullmq';
new Worker('exports', async (job) => {
const rows = await countRows(job.data.query);
let done = 0;
let lastReport = 0;
for await (const batch of readBatches(job.data.query, 1000)) {
await writeBatch(job.id, batch);
done += batch.length;
if (Date.now() - lastReport > 250 || done === rows) { // time-throttled
lastReport = Date.now();
await job.updateProgress({ done, total: rows, pct: Math.floor((done / rows) * 100), stage: 'rows' });
}
}
return { url: `/downloads/${job.id}.csv` }; // becomes the completed value
}, { connection });
Step 2 — Relay BullMQ events to SSE streams Permalink to this section
One QueueEvents instance per process receives every event for the queue; route them to the streams watching that job.
// relay.js — one QueueEvents per process, fan-out to local streams by job id.
import { QueueEvents, Job } from 'bullmq';
const events = new QueueEvents('exports', { connection });
const watchers = new Map(); // jobId → Set<res>
const send = (jobId, event, data, id) => {
for (const res of watchers.get(jobId) ?? []) {
res.write(`event: ${event}\n${id ? `id: ${id}\n` : ''}data: ${JSON.stringify(data)}\n\n`);
if (event === 'done' || event === 'failed') res.end();
}
};
events.on('progress', ({ jobId, data }) => send(jobId, 'progress', data, `p${data.done}`));
events.on('completed', ({ jobId, returnvalue }) => send(jobId, 'done', { state: 'succeeded', ...returnvalue }, 'end'));
events.on('failed', ({ jobId, failedReason }) => send(jobId, 'failed', { state: 'failed', message: failedReason }, 'end'));
Step 3 — Serve the stream, starting from the job’s current state Permalink to this section
app.get('/api/jobs/:id/events', requireJobAccess, async (req, res) => {
const job = await Job.fromId(exportsQueue, req.params.id);
if (!job) return res.status(404).end();
if (req.get('Last-Event-ID') === 'end') return res.status(204).end(); // already finished
openStream(res, { retryMs: 2000 });
let set = watchers.get(job.id);
if (!set) watchers.set(job.id, (set = new Set()));
set.add(res); // register before reading state
req.on('close', () => { set.delete(res); if (!set.size) watchers.delete(job.id); });
const state = await job.getState(); // waiting | active | completed | failed
if (state === 'completed') return send(job.id, 'done', { state: 'succeeded', ...job.returnvalue }, 'end');
if (state === 'failed') return send(job.id, 'failed', { state: 'failed', message: job.failedReason }, 'end');
res.write(`event: state\ndata: ${JSON.stringify({ state })}\n\n`);
if (job.progress?.done) res.write(`event: progress\ndata: ${JSON.stringify(job.progress)}\n\n`);
});
Registering the watcher before reading state closes the race where the job completes between the read and the registration. At worst, a progress update is sent twice, which the monotonic progress bar ignores.
Step 4 — Report and publish from a Celery task Permalink to this section
Celery has no push channel, so the task publishes its own frames next to update_state.
# tasks.py — Celery
import json, time, redis
from celery import shared_task
r = redis.Redis.from_url(REDIS_URL)
def report(task, event, data):
frame = json.dumps({"event": event, "data": data})
task.update_state(state="PROGRESS" if event == "progress" else event.upper(), meta=data)
r.publish(f"job:{task.request.id}", frame)
@shared_task(bind=True)
def export_rows(self, query):
total = count_rows(query)
done, last = 0, 0.0
try:
for batch in read_batches(query, 1000):
write_batch(self.request.id, batch)
done += len(batch)
if time.monotonic() - last > 0.25 or done == total:
last = time.monotonic()
report(self, "progress", {"done": done, "total": total, "pct": done * 100 // total})
result = {"state": "succeeded", "url": f"/downloads/{self.request.id}.csv"}
report(self, "done", result)
return result
except Exception as exc:
report(self, "failed", {"state": "failed", "message": str(exc)[:200]})
raise
Step 5 — Serve the Celery stream from FastAPI Permalink to this section
from celery.result import AsyncResult
@app.get("/api/jobs/{job_id}/events")
async def job_events(job_id: str, request: Request):
pubsub = aredis.pubsub()
await pubsub.subscribe(f"job:{job_id}") # subscribe first
res = AsyncResult(job_id)
async def gen():
try:
yield "retry: 2000\n\n"
if res.state == "SUCCESS":
yield f"event: done\nid: end\ndata: {json.dumps(res.result)}\n\n"; return
if res.state == "FAILURE":
yield f"event: failed\nid: end\ndata: {json.dumps({'state': 'failed'})}\n\n"; return
if res.state == "PROGRESS":
yield f"event: progress\ndata: {json.dumps(res.info)}\n\n"
while not await request.is_disconnected():
msg = await pubsub.get_message(ignore_subscribe_messages=True, timeout=15)
if msg is None:
yield ": hb\n\n"; continue
f = json.loads(msg["data"])
ident = "id: end\n" if f["event"] in ("done", "failed") else ""
yield f"event: {f['event']}\n{ident}data: {json.dumps(f['data'])}\n\n"
if ident:
return
finally:
await pubsub.aclose()
return StreamingResponse(gen(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})
Reading AsyncResult touches the result backend once per connection, not per second, which is the difference from the polling design.
Validation & Monitoring Permalink to this section
# Enqueue, stream, and confirm the stream terminates without a client-side close.
id=$(curl -s -X POST https://app.example.com/api/exports -d '{"query":"all"}' \
-H 'Content-Type: application/json' | jq -r .jobId)
time curl -sN https://app.example.com/api/jobs/$id/events | grep -E '^event:' | uniq -c
# 1 event: state
# 38 event: progress
# 1 event: done
# real 0m9.8s ← curl exited when the server ended the response
# Redis should show pub/sub traffic, not a stream of HGET polls.
redis-cli --latency-history -i 5 &
redis-cli MONITOR | grep -c HGETALL # should stay near zero while streams are open
Monitor stream duration against job duration. A stream that outlives its job by more than the retry interval means a terminal event is not ending the response, or the client is not closing on it.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Can I use BullMQ's QueueEvents directly in every request handler?
Create one instance per process and route by job id. Each QueueEvents instance holds its own blocking Redis connection, so one per request would exhaust connections quickly.
Does Celery's result backend need to be Redis?
Not for state, but the progress channel does need a pub/sub system. Using Redis for both keeps the task helper simple; with a database backend, publish to Redis or NATS separately.
What id should progress events use?
An increasing value such as the rows done, so the server can skip frames the client already has. Use a fixed id such as end for the terminal event so that a returning client can be recognised and told to stop.
How do I report progress for jobs with unknown total work?
Send stage names and counts instead of a percentage, and render an indeterminate progress bar with the stage text. A made-up percentage that stalls at 95 % erodes trust faster than an honest spinner.