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/:id every 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.

Where each queue exposes progress Matrix comparing BullMQ and Celery on where progress is stored, whether a push notification exists, and what the SSE handler should subscribe to. Where each queue exposes progress Queue Progress stored in Push available SSE handler listens to BullMQ job hash QueueEvents progress, completed, failed Celery result backend no own Redis channel Custom your store if you publish job channel
BullMQ already publishes progress events through Redis; Celery does not, so the worker publishes them itself.

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.

BullMQ progress reaching the browser Sequence diagram of a worker calling updateProgress, the queue writing an event, the relay's QueueEvents receiving it, and the SSE handler writing a progress frame; completion ends the stream. BullMQ progress reaching the browser Worker Redis Relay Browser updateProgress 25 % QueueEvents progress event: progress return value (completed) QueueEvents completed event: done, then end
No component polls. The worker's write is the trigger for every frame the browser receives.

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
Redis commands per second with 500 users watching jobs Bar chart comparing Redis command rates for browser polling, server-side polling inside the SSE handler, and event-driven relay. Redis commands per second with 500 users watching jobs Browser polls every 2 s 250 / s + 250 HTTP req/s SSE handler polls every 1 s 500 / s Event-driven relay ~12 / s (worker writes only) Redis commands per second
Moving the poll from the browser into the SSE handler saves HTTP overhead but not Redis load. Only the event-driven relay decouples load from viewer count.

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.