Dropping vs Coalescing Events Under Backpressure Permalink to this section
Part of Rate Limiting & Backpressure Handling, under Backend Stream Generation & Connection Management.
Every Server-Sent Events server eventually meets a client that cannot keep up: a phone on a weak signal, a browser tab starved of CPU, a corporate proxy that trickles bytes. At that moment the server must choose between four things — buffer more, drop events, merge events, or disconnect. Many servers never choose explicitly, which means they buffer until memory runs out. This guide lays out the options, shows which fits which data shape, and implements each so the choice is deliberate.
Symptom & Developer Intent Permalink to this section
- Server memory grows with the number of mobile users, not with total traffic.
- Some clients see prices or dashboard values that are minutes old while their connection is still open.
- Notifications occasionally go missing for users on poor connections, with nothing in the logs.
- Disconnecting slow clients causes them to reconnect and immediately fall behind again.
- Different event types on the same stream need different treatment, and the code has one policy for all.
The intent is a per-event-type policy chosen from the shape of the data, bounded memory per client, and correctness guarantees that are explicit rather than accidental.
Root Cause Analysis Permalink to this section
When a client reads slower than events are produced, the difference accumulates somewhere. Unbounded buffering stores it on the server, so memory grows and the client sees ever-older data. Every alternative bounds the buffer and decides what to give up when it is full.
The four policies:
| Policy | What happens when the client is behind | Correct for | Wrong for |
|---|---|---|---|
| Coalesce by key | new value replaces the queued one for the same key | prices, gauges, presence, progress | append-only feeds |
| Buffer then disconnect | bounded queue; on overflow, end the stream; client replays from its id | notifications, messages, audit feeds | high-rate state data |
| Drop when full | discard new events for this client | cursors, typing signals, animations | anything the client must see |
| Resnapshot | discard the queue, send a full snapshot when the client recovers | large dashboards, boards, documents | streams without a snapshot endpoint |
Step-by-Step Resolution Permalink to this section
Step 1 — Classify every event type Permalink to this section
const POLICY = {
price: { kind: 'coalesce', key: (d) => d.sym },
progress: { kind: 'coalesce', key: (d) => d.jobId },
notification: { kind: 'buffer', max: 256 },
message: { kind: 'buffer', max: 256 },
cursor: { kind: 'drop' },
typing: { kind: 'drop' },
};
Step 2 — Give each client a queue that applies the policy Permalink to this section
class ClientQueue {
constructor(res) {
this.res = res;
this.ordered = []; // buffer-policy frames, in order
this.latest = new Map(); // coalesce-policy frames, one per key
this.writing = false;
}
offer(type, data, frame) {
const p = POLICY[type] ?? { kind: 'buffer', max: 256 };
if (!this.writing) { this.write(frame); return; } // fast path: socket keeping up
switch (p.kind) {
case 'coalesce': this.latest.set(`${type}:${p.key(data)}`, frame); break; // replace, don't append
case 'drop': break; // ephemeral: skip
case 'buffer':
if (this.ordered.length >= p.max) return this.overflow(); // must not lose: reconnect
this.ordered.push(frame);
break;
}
}
write(frame) {
if (!this.res.write(frame)) {
this.writing = true;
this.res.once('drain', () => this.flush());
}
}
flush() {
this.writing = false;
const frames = [...this.ordered, ...this.latest.values()];
this.ordered = []; this.latest.clear();
for (const f of frames) {
if (this.writing) { this.ordered.push(f); continue; } // backed up again mid-flush
this.write(f);
}
}
overflow() {
this.res.write('retry: 1000\n\n');
this.res.end(); // client reconnects with Last-Event-ID and replays
}
}
The queue writes directly whenever the socket is keeping up, so fast clients never touch the policy code. Only a backed-up client accumulates anything, and what it accumulates is bounded: one entry per coalescing key, and at most max ordered frames.
Step 3 — Keep ordering guarantees honest Permalink to this section
Coalescing changes delivery order: a coalesced price may be written after a notification that was published later. That is acceptable when event types are independent. When they are not — a “position closed” notification that must follow the final price — put both in the same policy class, or attach the price to the notification.
Step 4 — Make buffer-then-disconnect converge Permalink to this section
A client disconnected for being slow reconnects and replays. If its connection is persistently too slow for the stream, it will fall behind and be disconnected again, in a loop. Break the loop by sending state-shaped data as a snapshot on reconnect instead of replaying every intermediate value, and by capping replay size — beyond the cap, send resync and let the client fetch current state over HTTP.
Step 5 — Measure what each policy is doing Permalink to this section
metrics.coalesced.inc({ type }); // replaced before being sent
metrics.dropped.inc({ type }); // intentionally discarded
metrics.slowDisconnects.inc(); // buffer overflow → reconnect
metrics.queueDepth.observe(this.ordered.length + this.latest.size);
Validation & Monitoring Permalink to this section
# Throttle a client to 2 KB/s against a busy stream and watch its behaviour.
curl -sN --limit-rate 2k https://app.example.com/api/stream > slow.log &
sleep 60
grep -c '^event: price' slow.log # far fewer than published, all recent
grep -c '^event: notification' slow.log # every notification, or a reconnect with replay
Dashboards should show, per event type, the ratio of coalesced and dropped events to delivered ones, and slow disconnects per minute. Spikes that align with a particular region or carrier are normal network conditions; spikes across all clients mean the server itself is too slow.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Why not just buffer everything and let slow clients catch up?
A client that is persistently slower than the stream never catches up. Its buffer grows until the server runs out of memory, and everything it does receive is increasingly stale.
Is dropping events ever safe?
For ephemeral signals such as cursor positions and typing indicators, yes — they are meaningless seconds later. For anything the client is expected to show or count, no.
Where does backpressure show up in Node.js?
res.write returns false when the socket's buffer passes its high-water mark, and res.writableNeedDrain stays true until the drain event. Those are the signals the queue uses to switch into policy mode.
Should coalescing happen on the server or the client?
Both. The server coalesces to save bandwidth for slow clients; the client coalesces to the display frame rate to save rendering work. They solve different bottlenecks.