Real-Time Application Patterns Permalink to this section
The other three sections of this site explain the parts: the event stream wire format, the server side that generates it, and the browser side that consumes it. This section is about the products you build from those parts. A live operations dashboard, a notification bell, a progress bar for a twenty-minute export, a build log that scrolls as the compiler runs, a list of who else is looking at a document, a price board — each is a Server-Sent Events stream underneath, and each one goes wrong in its own characteristic way. The audience is the engineer who has a working EventSource demo and now has to ship the feature: decide what goes in the stream, how much of it, how often, what happens when the connection drops, and what the server must remember so that the user never sees a gap.
The pages that follow are organised by the shape of the data rather than by technology, because the shape decides the design. A dashboard is a set of values that change; the only thing a late subscriber needs is the latest value. A notification feed is an append-only list; a late subscriber needs everything it missed. A job progress stream is a finite sequence with a known end. A log is an unbounded sequence where the reader, not the server, decides how far back to look. Presence is a set whose members expire. Market data is a firehose where most messages are obsolete before they are rendered. Get the shape right and the reconnect semantics, the replay buffer and the backpressure policy follow almost mechanically.
Concept Overview Permalink to this section
Every pattern in this section reduces to the same three-part pipeline: a source of change (a database write, a queue message, a worker reporting progress, a log line), a stream projector that turns changes into SSE frames for one subscriber, and a client reducer that folds frames into whatever the interface renders. The protocol is only the middle hop. The design work is deciding what the projector emits and what the reducer assumes.
The projector is where the per-feature decisions live. It filters (this user only sees their own notifications), shapes (a dashboard receives the delta, not the whole row), stamps an id so a reconnect can resume, and sometimes merges (ten price ticks in the same 100 ms become one). A minimal projector for a notification feed looks like this:
// One projector per open stream. It knows who the subscriber is and where they are up to.
async function* notificationProjector(userId, lastEventId, bus) {
// 1. Replay anything newer than the cursor the browser sent back on reconnect.
for (const n of await db.notificationsAfter(userId, lastEventId)) {
yield { id: n.id, event: 'notification', data: n };
}
// 2. Then follow the live bus, filtered to this subscriber.
for await (const n of bus.subscribe(`user:${userId}:notifications`)) {
yield { id: n.id, event: 'notification', data: n };
}
}
// The transport is generic: it only serialises whatever the projector yields.
app.get('/api/notifications/stream', async (req, res) => {
res.writeHead(200, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'X-Accel-Buffering': 'no',
});
const cursor = req.get('Last-Event-ID') ?? req.query.after ?? '0';
for await (const f of notificationProjector(req.user.id, cursor, bus)) {
res.write(`id: ${f.id}\nevent: ${f.event}\ndata: ${JSON.stringify(f.data)}\n\n`);
}
});
The split matters because it keeps the transport identical across every feature. Heartbeats, handling client disconnects, headers and flush behaviour are written once. What changes from feature to feature is only the projector.
The client reducer is the mirror image. It receives frames in order, applies each one to local state, and hands that state to the rendering layer. Keeping the reducer pure — a function of the previous state and one frame — has a practical benefit that becomes obvious the first time a bug report arrives: a recorded stream can be replayed through the reducer in a unit test and the exact interface state reproduced without a server.
// A reducer for a dashboard stream: snapshot replaces, delta merges, anything else is ignored.
export function dashboardReducer(state, frame) {
switch (frame.type) {
case 'snapshot': return { version: frame.id, values: frame.data };
case 'delta': return { version: frame.id, values: { ...state.values, ...frame.data } };
default: return state; // unknown events must never break the view
}
}
The six shapes, one at a time Permalink to this section
Current state (dashboards). A set of named values, each of which is only interesting at its latest version. The stream opens with a snapshot and carries deltas afterwards; a reconnect simply opens with a fresh snapshot. Nothing needs to be buffered on the server, which is why dashboards scale further than any other pattern per node. The traps are on the client — rendering on every frame, and charts that grow without bound. The live dashboards and metrics feeds topic covers both.
Append-only lists (notifications and activity). Every item matters and order matters. The stream must be resumable, which means event ids that the server can query by, a replay window long enough to cover a laptop lid being closed over lunch, and a client that tolerates the occasional duplicate. Unread counts are a separate, state-shaped value that rides the same connection. See notification and activity feeds.
Finite sequences (job progress). The stream has a beginning, a known end and usually a terminal event — done or failed — after which the server closes the response. The client must be able to arrive late (the user opens the page after the job started) and must not reconnect forever after the job has finished. Progress streaming for long-running jobs walks through the state machine.
Unbounded sequences (logs and build output). The volume is high, bursty and mostly uninteresting, and the reader decides how much history they want. The server streams from an offset, bounds what it will replay, and applies filters before bytes leave the process. Log tailing and CI output streaming covers offsets, ANSI colour and server-side filtering.
Expiring sets (presence). Membership is asserted by being connected and revoked by silence. Treating presence as a lease rather than a join/leave log is the whole trick. Collaborative presence and live updates pairs presence with the POST-up, SSE-down pattern for shared editing.
Firehoses (market data and scores). The upstream rate exceeds what a human can read and often what a phone can render. Intermediate values are disposable, so the server conflates — keeps only the latest value per key and flushes on a timer — and slow clients simply see fewer, fresher updates. Market data and live scoreboards shows conflation and sequence gap detection.
Specification & Wire Format Permalink to this section
The wire format does not change between patterns, but how each pattern uses the four fields does. The single most useful design artefact for a new real-time feature is a table like the one below, filled in before any code is written.
| Pattern | event: names |
id: meaning |
Replay on reconnect | Typical rate |
|---|---|---|---|---|
| Live dashboard | snapshot, delta |
version of the whole view | none — send a fresh snapshot | 1–10 per second |
| Notification feed | notification, read |
monotonic per user | everything after the cursor | a few per hour |
| Job progress | progress, done, failed |
job id + sequence | current state only | 1–5 per second, finite |
| Log tail | line, truncated |
byte offset or line number | from the offset, bounded | bursts of thousands |
| Presence | join, leave, roster |
none, or roster version | full roster | on change + expiry sweeps |
| Market data | tick, book |
exchange sequence | latest value per symbol | hundreds per second |
Two wire-level conventions pay for themselves in every pattern. First, name every event. An unnamed stream forces every message through onmessage and a type switch in JSON, which works until two features share a connection. Named event types let each consumer attach its own listener. Second, make the id mean something the server can seek to. An opaque UUID is useless as a replay cursor; a per-stream sequence number, a database row id or a byte offset is exactly what the replay query needs.
: dashboard stream — every frame names its type and carries a version the server can compare
event: snapshot
id: v1842
data: {"cpu":41,"rps":1290,"errors":3,"p95":212}
event: delta
id: v1843
data: {"rps":1302}
event: delta
id: v1844
data: {"errors":4,"p95":230}
Architecture Patterns Permalink to this section
Three deployment topologies cover almost every feature in this section. The choice between them is driven by how many subscribers share the same data.
Shared-topic fan-out. Many subscribers, identical data: dashboards, scoreboards, a public status page. Publish once to a broker and let each SSE node multiply locally. The expensive failure is a per-connection database query triggered by every change, which turns 10,000 viewers into 10,000 queries per update. Redis pub/sub fan-out is the usual first implementation.
Per-subscriber channels. Each subscriber sees different data: notifications, a user’s own jobs, per-account activity. Publish to a channel keyed by the user or tenant. The failure here is the opposite one — a broadcast channel with client-side filtering, which leaks other users’ events to anyone who opens DevTools.
Owned streams. One producer, one consumer, a finite life: an upload being processed, an export, an AI generation. The stream is keyed by the job id, often lives on the node that runs the job, and ends when the job ends. Its failure mode is a load balancer that routes the reconnect to a different node that has never heard of the job.
# Owned streams: route the progress stream for a job back to the node that owns it.
# The job id is part of the path, so a consistent hash keeps reconnects on the same node.
upstream job_nodes {
hash $job_id consistent;
server 10.0.1.11:8080;
server 10.0.1.12:8080;
}
map $uri $job_id {
~^/api/jobs/(?<id>[^/]+)/events$ $id;
}
location ~ ^/api/jobs/[^/]+/events$ {
proxy_pass http://job_nodes;
proxy_http_version 1.1;
proxy_set_header Connection "";
proxy_buffering off;
proxy_read_timeout 1h;
}
Authorisation belongs in the projector Permalink to this section
Whichever topology you choose, authorisation must be enforced where frames are produced, not where they are rendered. A shared-topic stream must still check that the subscriber may see the topic at connect time, and a per-subscriber stream must derive the channel name from the authenticated identity rather than from a query parameter. /api/notifications/stream?user=42 is a data leak waiting for someone to change the number. The channel is user:${req.user.id}, always, and the authentication options for SSE streams decide how req.user is established on a request that EventSource cannot add headers to.
Long-lived connections add one more wrinkle: permissions can change while the stream is open. A user removed from a project should stop receiving its events within seconds, not at their next reconnect. The simplest reliable approach is to publish a revoke message on the user’s control channel that makes the node close that user’s streams; the client reconnects, and the connect-time check does the rest.
The sticky route is a convenience, not a correctness guarantee — nodes restart. The robust version stores job state somewhere every node can read, so that any node can answer a reconnect with the current state. Resuming progress streams after a reconnect shows both halves.
Edge Cases & Failure Modes Permalink to this section
The failures below recur across every pattern. Each has a specific symptom that points straight at it.
Stale state after reconnect.
- Symptom: after a network blip the dashboard keeps showing the old numbers until something changes.
- Root cause: the stream only sends deltas; a reconnecting client has no baseline.
- Mitigation: send a
snapshotevent as the first frame of every connection, then deltas.
Duplicates after replay.
- Symptom: the same notification appears twice, usually right after the laptop wakes.
- Root cause: the client rendered an event, the connection dropped before its
idwas recorded, and the server replayed it. - Mitigation: key the client store by event id; see deduplicating events on the client.
Ghost presence.
- Symptom: users stay “online” for hours after they left.
- Root cause: presence relies on a
leaveevent, and a connection that dies silently never sends one. - Mitigation: presence is a lease renewed by the connection’s heartbeat; absence of renewal is the leave.
Render storms.
- Symptom: the tab becomes unresponsive when the feed gets busy.
- Root cause: every message triggers a synchronous re-render.
- Mitigation: accumulate into a buffer and flush once per animation frame, as in throttling dashboard updates to the frame rate.
Unbounded client memory.
- Symptom: a log viewer left open overnight uses gigabytes.
- Root cause: every line is appended to an array that is never trimmed.
- Mitigation: cap the client buffer and virtualise the list; the server keeps the full history.
Horizontal Scaling & Production Ops Permalink to this section
Real-time features change the capacity model of an application. A request/response service is sized by requests per second; a streaming service is sized by concurrent connections multiplied by message rate. The two dimensions stress different resources — connections cost memory and file descriptors, messages cost CPU and bandwidth — and each pattern sits in a different corner.
The operational rules that follow from that model:
- Budget file descriptors per node, not per service. Each open stream holds one socket; raise
ulimit -nand the kernel limits together, as covered in tuning file descriptor limits. - Separate the firehose. Host market data or log tails on a different process pool from notifications, so a busy feed cannot starve a quiet but important one.
- Measure delivery latency, not only throughput. The metric users feel is the time from source change to render. Stamp events at the source and compare at the edge, as described in measuring end-to-end event latency.
- Drain before deploy. Every deploy disconnects every stream. Stagger the restart and send a
retry:hint with jitter so reconnects spread out instead of arriving as one wave.
Observability for these features has to answer a question that request metrics never ask: is the user looking at the truth right now? Three signals together answer it. Open streams per feature (a gauge, labelled by event family rather than by URL) shows whether clients are connected at all. Frames written per second per feature shows whether the projector is producing. Source-to-write lag — the difference between the timestamp stamped on the change and the moment the frame hits the socket — shows whether what is being produced is current. A dashboard that is connected, busy and ninety seconds behind is broken in a way no error-rate alert will catch, and it is the most common failure of a fan-out tier under load. The observability and metrics topic has the instrumentation.
// Stamp at the source; measure at the write. The difference is what users experience.
bus.on('change', (evt) => { evt.sourceTs ??= Date.now(); });
function writeFrame(res, feature, frame) {
res.write(`id: ${frame.id}\nevent: ${frame.event}\ndata: ${JSON.stringify(frame.data)}\n\n`);
metrics.framesWritten.inc({ feature });
metrics.sourceToWriteMs.observe({ feature }, Date.now() - frame.sourceTs);
}
# Kubernetes: give streams time to hand over before the pod is killed.
spec:
terminationGracePeriodSeconds: 60
containers:
- name: sse
lifecycle:
preStop:
exec:
# Stop accepting new streams, tell existing ones to reconnect elsewhere.
command: ["/bin/sh", "-c", "curl -s -X POST localhost:8080/admin/drain && sleep 45"]
readinessProbe:
httpGet: { path: /ready, port: 8080 }
periodSeconds: 2
Migration & Fallback Paths Permalink to this section
Most real-time features start life as polling. The migration to SSE is almost always incremental, and the safest order keeps the polling endpoint as the fallback for the whole migration.
- Shadow. Open the stream alongside the existing poll and compare what each delivers. Log divergence rather than rendering from the stream.
- Switch with fallback. Render from the stream; if it fails to open within a few seconds, or errors repeatedly, drop back to polling. The client code for this degradation is short:
function subscribe(url, onData, { pollUrl, pollMs = 5000 } = {}) {
let es, timer, failures = 0;
const startPolling = () => {
if (timer) return;
const tick = async () => onData(await (await fetch(pollUrl)).json());
tick();
timer = setInterval(tick, pollMs);
};
es = new EventSource(url);
const opened = setTimeout(startPolling, 4000); // no open event in 4 s → degrade
es.onopen = () => { clearTimeout(opened); failures = 0; clearInterval(timer); timer = null; };
es.addEventListener('snapshot', (e) => onData(JSON.parse(e.data)));
es.onerror = () => { if (++failures >= 3) startPolling(); };
return () => { es.close(); clearTimeout(opened); clearInterval(timer); };
}
- Retire polling once the stream has survived a deploy, a proxy change and a traffic spike. Keep the snapshot endpoint — the stream’s first frame and the fallback share it.
Where EventSource itself is the obstacle — a POST body, custom headers, or a runtime without it — a fetch-based SSE client reads the same wire format through a ReadableStream. And when the question is whether SSE is the right transport at all, SSE vs WebSockets vs HTTP polling has the decision matrix.
⚡ Production Directives
- Classify every real-time feature by data shape — current state, append-only, finite, unbounded or expiring — before choosing a replay strategy.
- Send a full snapshot as the first frame of every state-shaped stream; never make a reconnecting client wait for the next change.
- Use an
idthe server can seek to: a sequence, a row id or a byte offset, never a random UUID. - Batch client rendering to the animation frame for any feed that can exceed ten messages per second.
- Keep the polling path as a tested fallback until the stream has survived a real incident.
Frequently Asked Questions Permalink to this section
Can several real-time features share one SSE connection?
Yes, and on HTTP/1.1 they usually should, because browsers cap connections per origin at six. Give each feature its own named event type, multiplex them on one stream, and let each component add its own listener. On HTTP/2 separate streams are cheap and simpler to reason about.
Should a dashboard stream send deltas or full values?
Both: a full snapshot as the first frame of each connection, then deltas. Deltas keep bandwidth low while connected, and the snapshot makes reconnects correct without a replay buffer.
How do I know which replay strategy a feature needs?
Ask what a user who was disconnected for thirty seconds must see when they come back. If it is every item they missed, replay after the cursor. If it is only the current value, send a snapshot. If it is the tail of something long, replay from an offset with a bound.
How long should the server keep events for replay?
Long enough to cover the disconnections your users really have. Mobile clients commonly drop for tens of seconds and laptops for hours. A notification feed backed by a database table can replay indefinitely; an in-memory buffer should cover at least a few minutes, and anything older should trigger a resync signal rather than a silent gap.
What is the biggest capacity risk when adding a real-time feature?
A per-connection database query on every change. It turns viewer count into query load and fails suddenly at a traffic spike. Publish each change once, fan it out in memory on every node, and reserve database reads for the snapshot at connect time.
Do real-time features need a separate service?
Not at first. A single process can hold tens of thousands of idle streams. Split the streaming tier out when its deploy cadence, scaling profile or failure domain needs to differ from the request/response API — typically when a firehose feed arrives.
Is SSE suitable for collaborative editing?
For the server-to-client half, yes: SSE carries remote changes, presence and cursors well, while edits go up as ordinary POST requests. WebSockets only become necessary when the upstream rate is high enough that per-request overhead dominates.