Progress Streaming for Long-Running Jobs Permalink to this section
Part of Real-Time Application Patterns.
Exports, data imports, video transcodes, report generation, deployments, large AI tasks — any operation that takes longer than a user will wait on a spinner needs a progress stream. The request that starts the job returns immediately with a job id; the browser then opens a Server-Sent Events stream for that id and watches a progress bar fill. The pattern looks simple and hides three hard problems: the worker that does the job is usually not the process that holds the stream, the user can arrive late or reconnect at any point, and the stream has to end — a progress stream that reconnects forever after the job finished is a slow resource leak. This guide covers the job state model, how progress travels from worker to browser, the terminal events that close the loop, and the client that handles all of it. It is for backend and full-stack engineers adding asynchronous work to a product, and the code is shown in Node.js and Python.
How It Works Permalink to this section
A job is a small state machine, and the progress stream is a projection of it. Every state transition, and every progress update within the running state, becomes an event.
On the wire the stream is short and ends deliberately:
retry: 2000
event: state
data: {"jobId":"exp_7Qk","state":"queued","position":3}
event: state
data: {"jobId":"exp_7Qk","state":"running","total":48210}
event: progress
id: 12
data: {"done":12050,"total":48210,"pct":25,"stage":"rows","etaS":41}
event: progress
id: 13
data: {"done":24100,"total":48210,"pct":50,"stage":"rows","etaS":27}
event: done
id: 99
data: {"state":"succeeded","url":"/downloads/exp_7Qk.csv","bytes":8839210}
The terminal event is followed by the server ending the response. Here lies the first subtlety: when the response ends, EventSource treats it as a dropped connection and reconnects. The client must call close() when it receives a terminal event, and the server should answer any later request for a finished job with the terminal event again — or with HTTP 204, which tells EventSource to stop reconnecting for good. The retry rules are covered in setting the retry interval in SSE streams.
Progress is state-shaped, not a log. A client that reconnects during the job needs the current percentage and stage, not the forty intermediate updates it missed. So the stream opens with the current state as its first frame, and id values serve to skip that frame when the client is already current, not to drive replay.
Server-Side Implementation Permalink to this section
The worker and the stream handler rarely share a process. The worker publishes progress to a store every node can read and a channel every node can subscribe to; the stream handler reads the current state on connect and then follows the channel.
The worker side reports progress through a small helper that both persists and publishes, and that rate-limits itself so a fast loop does not flood the channel:
// worker/progress.js
export function progressReporter(jobId, { minIntervalMs = 250 } = {}) {
let last = 0;
let seq = 0;
const key = `job:${jobId}`;
async function emit(event, data) {
seq += 1;
const frame = { event, id: seq, data };
await redis.multi()
.hset(key, { state: data.state ?? 'running', last: JSON.stringify(frame), updatedAt: Date.now() })
.expire(key, 24 * 3600) // finished jobs are kept for a day
.publish(key, JSON.stringify(frame))
.exec();
}
return {
async progress(done, total, stage) {
const now = Date.now();
if (now - last < minIntervalMs && done < total) return; // throttle, but never drop 100 %
last = now;
await emit('progress', { done, total, pct: Math.floor((done / total) * 100), stage });
},
succeeded: (result) => emit('done', { state: 'succeeded', ...result }),
failed: (code, message) => emit('failed', { state: 'failed', code, message }),
cancelled: () => emit('cancelled', { state: 'cancelled' }),
};
}
The stream handler reads the latest frame from the hash, sends it, and either closes immediately (terminal state) or follows the channel until a terminal frame arrives:
// api/job-events.js
const TERMINAL = new Set(['done', 'failed', 'cancelled']);
app.get('/api/jobs/:id/events', requireJobAccess, async (req, res) => {
const key = `job:${req.params.id}`;
const snap = await redis.hgetall(key);
if (!snap.last) return res.status(404).end();
const latest = JSON.parse(snap.last);
if (TERMINAL.has(latest.event) && Number(req.get('Last-Event-ID')) === latest.id) {
return res.status(204).end(); // client already has the ending: stop retrying
}
openStream(res, { retryMs: 2000 });
const write = (f) => res.write(`event: ${f.event}\nid: ${f.id}\ndata: ${JSON.stringify(f.data)}\n\n`);
const sub = redis.duplicate();
await sub.subscribe(key); // subscribe before sending state: no gap
write(latest);
if (TERMINAL.has(latest.event)) return end();
sub.on('message', (_c, raw) => {
const f = JSON.parse(raw);
if (f.id <= latest.id) return; // already covered by the snapshot
write(f);
if (TERMINAL.has(f.event)) end();
});
req.on('close', end);
function end() { sub.quit().catch(() => {}); res.end(); }
});
The equivalent in Python keeps the same order — subscribe, read state, send, follow — with redis.asyncio:
TERMINAL = {"done", "failed", "cancelled"}
@app.get("/api/jobs/{job_id}/events")
async def job_events(job_id: str, request: Request, user=Depends(job_access)):
key = f"job:{job_id}"
pubsub = redis.pubsub()
await pubsub.subscribe(key)
snap = await redis.hgetall(key)
if not snap:
await pubsub.aclose()
raise HTTPException(404)
latest = json.loads(snap[b"last"])
async def gen():
try:
yield "retry: 2000\n\n"
yield frame(latest)
if latest["event"] in TERMINAL:
return
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"])
if f["id"] <= latest["id"]:
continue
yield frame(f)
if f["event"] in TERMINAL:
return
finally:
await pubsub.aclose()
return StreamingResponse(gen(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})
def frame(f):
return f"event: {f['event']}\nid: {f['id']}\ndata: {json.dumps(f['data'])}\n\n"
Queue position while the job waits Permalink to this section
A job that sits in a queue for thirty seconds with a bar stuck at 0 % looks broken. While the job is queued, publish its position whenever it changes. The cheap way is to have the dispatcher, not every waiting job, compute positions: each time a worker takes a job, the dispatcher reads the next few waiting job ids and publishes a state frame with the new position for each of them. Users further back than, say, fifty places get a coarse message (“waiting for capacity”) rather than an exact number, which saves work and avoids promising precision the queue cannot keep.
// dispatcher: after a worker takes a job, update the positions of the next 50 waiters.
async function publishPositions(queueName) {
const waiting = await queue.getWaiting(0, 49); // oldest first
const pipe = redis.pipeline();
waiting.forEach((job, i) => {
const frame = { event: 'state', id: 0, data: { state: 'queued', position: i + 1 } };
pipe.hset(`job:${job.id}`, 'last', JSON.stringify(frame)).publish(`job:${job.id}`, JSON.stringify(frame));
});
await pipe.exec();
}
Queue frames use id: 0 so the first real progress frame (id 1) always supersedes them in the stream handler’s comparison.
Cancellation from the stream’s page Permalink to this section
Users will want to cancel a job they are watching. Cancellation is an ordinary POST — SSE is one-way — but its result must travel back through the stream as a terminal cancelled event, so that every tab watching the job ends in the same state:
app.post('/api/jobs/:id/cancel', requireJobAccess, async (req, res) => {
await redis.set(`job:${req.params.id}:cancel`, '1', 'EX', 3600); // the worker polls this flag
res.status(202).end();
});
// worker loop: check the flag between batches, never mid-write.
for (const batch of batches) {
if (await redis.exists(`job:${jobId}:cancel`)) {
await cleanupPartialOutput(jobId);
return reporter.cancelled(); // emits event: cancelled, a terminal state
}
await processBatch(batch);
await reporter.progress(done += batch.length, total, 'rows');
}
The 202 response says only that the request was accepted. The progress bar should switch to “cancelling…” on the click and to “cancelled” when the terminal event arrives, because a worker in the middle of a large batch may take several seconds to reach its next checkpoint.
Streaming background job progress with SSE wires this into a real queue (BullMQ and Celery), including the queued-position updates.
Client-Side Consumption Permalink to this section
The client opens the stream once it has a job id, renders each state, and closes on any terminal event.
// job-progress.js
export function watchJob(jobId, { onProgress, onDone, onFailed, onState }) {
const es = new EventSource(`/api/jobs/${jobId}/events`, { withCredentials: true });
let finished = false;
const finish = (fn, data) => { finished = true; es.close(); fn(data); };
es.addEventListener('state', (e) => onState?.(JSON.parse(e.data)));
es.addEventListener('progress', (e) => onProgress(JSON.parse(e.data)));
es.addEventListener('done', (e) => finish(onDone, JSON.parse(e.data)));
es.addEventListener('failed', (e) => finish(onFailed, JSON.parse(e.data)));
es.addEventListener('cancelled', (e) => finish(onFailed, { code: 'cancelled', ...JSON.parse(e.data) }));
es.onerror = () => {
if (!finished && es.readyState === EventSource.CLOSED) {
onFailed({ code: 'stream_closed', message: 'Lost contact with the job' });
}
};
return () => es.close();
}
Store the job id somewhere that survives navigation — the URL, or sessionStorage — so that a reload or returning to the page reopens the stream and immediately shows the current state. For progress bars, render pct with a native <progress> element, which is accessible without extra work, and show the stage text beside it. Two small rendering rules make a progress bar feel trustworthy. First, never let it move backwards: if a reconnect delivers a snapshot that is older than what is on screen — possible when the stream reconnects to a node whose store read raced a newer publish — keep the higher value. Second, animate between values with a short CSS transition rather than jumping, so 250 ms updates read as continuous motion:
let shown = 0;
function renderProgress(bar, label, { pct, stage, etaS }) {
shown = Math.max(shown, pct); // monotonic on screen
bar.value = shown; // <progress max="100">
label.textContent = etaS != null ? `${stage} — about ${formatEta(etaS)} left` : stage;
}
An ETA is useful only if it is stable: smooth it on the server over the last few updates, as covered in reporting file upload and processing progress.
Edge Cases & Network Interference Permalink to this section
- Stream before job. If the client opens the stream before the job record is written, the handler finds nothing. Write the job record in the same request that returns the id, before responding.
- Proxy buffering. A buffering proxy delivers progress in chunks, so the bar jumps from 10 % to 60 %. Disable buffering for the route; see disabling CDN buffering for event streams.
- Idle during long stages. A stage that takes ten minutes without progress updates (a database commit, an upload to object storage) will be closed by idle timeouts. Heartbeat comments cover it, and a periodic
stateframe withstagetext reassures the user. - Endless reconnect. A server that simply ends the response after
doneleavesEventSourcereconnecting everyretryinterval. The client mustclose(); the server should return 204 to a client that already has the terminal event. - Expired results. A download URL signed for fifteen minutes is useless when the user reopens the tab the next day. Resolve the result URL at render time from the job id, not from the event.
- Serverless timeouts. A function with a 30-second limit cannot hold a ten-minute progress stream. Either reconnect deliberately at the limit — the snapshot-on-connect design makes that seamless — or host the stream elsewhere; see handling execution timeouts in serverless SSE.
Mitigation checklist:
Performance & Scale Considerations Permalink to this section
Progress streams are short-lived and rarely numerous, but they have one sharp edge: a naive worker reports progress in its inner loop, and a job that processes a million rows publishes a million updates.
Time-based throttling (at most one update every 250 ms) gives smoother bars than percentage steps and bounds cost by duration. The reporter above does both: it throttles by time, but always sends the final 100 %.
Each open progress stream holds one broker subscription. At hundreds of concurrent jobs that is negligible; at tens of thousands — a batch import feature used by every customer at month end — move to node-addressed routing as in per-user channels for notification delivery. Keep terminal job records with a TTL (a day is typical) so late readers can still see the outcome without the store growing without bound.
Validation & Debugging Permalink to this section
# Start a job, then watch its stream end on its own.
job=$(curl -s -X POST -b s.txt https://app.example.com/api/exports | jq -r .jobId)
curl -sN -b s.txt https://app.example.com/api/jobs/$job/events
# Expect: state(queued) → state(running) → progress… → done, then curl exits.
# A late joiner sees the current state first.
curl -sN -b s.txt https://app.example.com/api/jobs/$job/events | head -3
# A client that already has the terminal event is told to stop.
curl -s -o /dev/null -w '%{http_code}\n' -b s.txt -H 'Last-Event-ID: 99' \
https://app.example.com/api/jobs/$job/events
# 204
In DevTools, the stream request should show a finite duration and a closed status after done, and no further requests to the same URL. A column of repeated requests every two seconds after completion is the endless-reconnect bug. Structured logs should connect the job and the stream:
{"evt":"job_stream_open","job":"exp_7Qk","state":"running","pct":50,"late_join":true}
{"evt":"job_stream_close","job":"exp_7Qk","reason":"terminal","frames":7,"duration_s":31}
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Why not just poll a job status endpoint?
Polling works and is a reasonable fallback. SSE gives immediate updates at lower cost for jobs that run seconds to minutes and removes the tuning of a polling interval. The same state endpoint can back both.
How do I stop EventSource reconnecting after the job finishes?
Call close() on the client when a terminal event arrives. As a server-side safety net, respond with HTTP 204 to a request whose Last-Event-ID already names the terminal event; the specification makes 204 stop reconnection permanently.
Should progress events carry an id?
Yes, a sequence per job. It lets the server skip frames the client already has and recognise a client that has seen the end. Replaying every intermediate progress update is not needed; the current state is.
How often should a worker report progress?
At most four times a second is plenty for a visible bar, and once a second is fine for jobs longer than a minute. Throttle by time rather than by work done so that large jobs do not publish more messages than small ones.
What if the worker crashes mid-job?
The stream stops receiving updates, and the store still says running. Have workers heartbeat the job record and a reaper mark jobs failed when the heartbeat is older than a threshold; the reaper publishes the failed event so watching clients end cleanly instead of waiting forever.
Can one stream watch several jobs?
Yes. For a jobs panel, open one stream for the user and send events tagged with job ids, fed by the user's channel. Per-job streams are simpler when the page shows a single job.