Replaying Missed Events with Redis Streams Permalink to this section
Part of Redis Pub/Sub Fan-Out for SSE, under Backend Stream Generation & Connection Management.
Redis pub/sub is fire-and-forget: a message published while a client is reconnecting is simply gone. For dashboards that is fine; for notifications, chat or order updates it is data loss. Redis Streams keep an ordered, trimmed log with a server-assigned id per entry — which is exactly what Last-Event-ID needs. This guide moves an SSE service from pub/sub to Streams, or adds Streams alongside pub/sub, so that every reconnect replays exactly what the client missed.
Symptom & Developer Intent Permalink to this section
- Users report missing notifications after their phone changes networks or their laptop wakes.
- The server ignores
Last-Event-IDbecause there is nothing to replay from. - A custom replay buffer in each node’s memory is lost on every deploy.
- Nodes disagree about event ids, so a client reconnecting to a different node gets the wrong replay.
The intent is a shared, durable, bounded log from which any node can replay any client’s gap, with ids that mean the same thing on every node.
Root Cause Analysis Permalink to this section
Replay needs three properties: a history that outlives connections, a position every node understands, and a way to seek to that position. Pub/sub has none of them. Per-node memory buffers have the first only until a restart and the second only if every node receives identical ids. A Redis Stream has all three: entries persist (up to a trimming limit), each has a monotonically increasing id like 1726650042000-0 assigned by Redis, and XRANGE returns entries after a given id.
The exclusive-range syntax (id (Redis 6.2+) returns entries strictly after the given id, which avoids resending the one the client already has.
Step-by-Step Resolution Permalink to this section
Step 1 — Publish with XADD and a trimming limit Permalink to this section
// Publisher: append to a per-audience stream, keeping roughly the last 10,000 entries.
const id = await redis.xadd(
`events:user:${userId}`, 'MAXLEN', '~', 10000, '*', // '*' = Redis assigns the id
'type', 'notification', 'data', JSON.stringify(payload));
MAXLEN ~ trims approximately, which is much cheaper than exact trimming. Choose the length from the replay window you promise: events per hour for the busiest audience times the hours of absence you want to repair.
Step 2 — Replay after the client’s id, then tail Permalink to this section
app.get('/api/stream', requireUser, async (req, res) => {
const key = `events:user:${req.user.id}`;
openStream(res, { retryMs: 3000 });
let last = req.get('Last-Event-ID') || '$'; // '$' = only new entries
if (last !== '$') {
const [[oldestId] = []] = await redis.xrange(key, '-', '+', 'COUNT', 1);
if (oldestId && compareIds(last, oldestId) < 0) {
res.write('event: resync\ndata: {"reason":"trimmed"}\n\n'); // gap exceeds retention
last = '$';
} else {
const entries = await redis.xrange(key, `(${last}`, '+', 'COUNT', 500);
for (const [id, fields] of entries) { write(res, id, fields); last = id; }
}
}
const reader = redis.duplicate(); // XREAD BLOCK needs its own connection
req.on('close', () => reader.disconnect());
while (!res.writableEnded) {
const out = await reader.xread('BLOCK', 15000, 'COUNT', 100, 'STREAMS', key, last).catch(() => null);
if (!out) { res.write(': hb\n\n'); continue; } // timeout doubles as heartbeat
for (const [id, fields] of out[0][1]) { write(res, id, fields); last = id; }
}
});
function write(res, id, fields) {
const f = Object.fromEntries(chunkPairs(fields));
res.write(`id: ${id}\nevent: ${f.type}\ndata: ${f.data}\n\n`);
}
$ resolves to “the last entry at the moment of the call”, so a new client starts from now. After the first read, last is always a concrete id, which makes the tail gap-free.
Step 3 — Scale beyond a connection per client Permalink to this section
A blocking XREAD per open stream means one Redis connection per viewer, which does not scale past a few thousand. The fix is the same as for pub/sub: one reader per node, fanning out in memory. The node reads many streams in a single XREAD call or reads a shared stream, and routes entries to local connections:
// One loop per node: tail the streams of users connected here, route locally.
async function tailLoop() {
for (;;) {
const keys = [...localUsers.keys()].map((u) => `events:user:${u}`);
if (!keys.length) { await sleep(200); continue; }
const ids = keys.map((k) => cursor.get(k) ?? '$');
const out = await reader.xread('BLOCK', 5000, 'COUNT', 500, 'STREAMS', ...keys, ...ids);
for (const [key, entries] of out ?? []) {
for (const [id, fields] of entries) { routeToLocal(key, id, fields); cursor.set(key, id); }
}
}
}
For very many users per node, publish to a single shared stream (or a handful of sharded ones) and filter in memory, keeping per-user streams only for replay. Scaling SSE across multiple nodes with Redis discusses the node-level fan-out.
Step 4 — Compare stream ids correctly Permalink to this section
Stream ids are milliseconds-sequence; compare the two numeric parts, not the strings:
function compareIds(a, b) {
const [am, as] = a.split('-').map(BigInt), [bm, bs] = b.split('-').map(BigInt);
return am === bm ? (as < bs ? -1 : as > bs ? 1 : 0) : (am < bm ? -1 : 1);
}
String comparison fails as soon as the millisecond parts have different lengths or the sequence part reaches two digits.
Validation & Monitoring Permalink to this section
# Inspect a user's stream and its bounds.
redis-cli XINFO STREAM events:user:42 | egrep 'length|first-entry|last-generated-id' -A1
# Replay: reconnect with an older id and confirm the next ids are strictly greater.
curl -sN -H 'Last-Event-ID: 1726650042000-3' https://app.example.com/api/stream | grep -m5 '^id:'
Monitor replay counts and sizes, resyncs due to trimming (the share of reconnects whose id had already been trimmed), stream lengths, and Redis memory. A rising trimmed-resync rate says the retention window is shorter than real client absences.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Should I replace pub/sub with Streams entirely?
Not necessarily. Many services keep pub/sub for live delivery and add Streams only as the replay store. Streams alone also work, with one reader per node doing the fan-out.
What about consumer groups?
Consumer groups distribute entries among workers for processing. SSE fan-out needs every node to see every entry, so plain XREAD is the right tool; consumer groups are for background processing of the same stream.
Can the SSE id be my own id instead of Redis's?
You can XADD with an explicit id, but it must increase monotonically per stream. Letting Redis assign ids avoids the coordination problem entirely.
How long can a client be away and still replay?
As long as its last id is still in the stream. MAXLEN and the publish rate together determine that window; measure it and publish it as part of the API contract.