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.

Choosing a backpressure policy Decision tree that picks a policy from the data shape: coalesce by key for current-state data, buffer then disconnect and replay for append-only data, drop for ephemeral signals, and resnapshot for large views. Choosing a backpressure policy Does only the latest value per key matter? Coalesce by key yes no Must every event be delivered? Buffer, then disconnect + replay yes no Is the event useless seconds later? Drop when full yes no Resnapshot on recovery
The data shape, not the transport, decides. One stream often carries several shapes and needs a policy per event type.

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.

How one slow client's queue handles a burst Flow of a burst of mixed events into a slow client's queue, splitting into coalesced prices, buffered notifications and dropped cursor signals, then a single flush when the socket drains. How one slow client's queue handles a burst Burst 200 mixed events offer Classify by event type price Coalesce 8 price keys notification Buffer 3 notifications flush On drain 11 frames 189 cursor and superseded price events were intentionally not sent
Two hundred incoming events became eleven outgoing frames, and none of the lost ones mattered: the prices were superseded and the cursors were stale.

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
Server memory per slow client after ten minutes Bar chart comparing memory held for one slow client under unbounded buffering, buffer-then-disconnect, and per-key coalescing on a busy mixed stream. Server memory per slow client after ten minutes Unbounded buffer ~10.8 MB, growing Buffer 256 + disconnect ≤ 90 KB Coalesce by key ~12 KB kilobytes held for one client reading at 2 KB/s from a 20 KB/s stream
Unbounded buffering stores the stream's history for the client. The bounded policies store at most a queue or one value per key.

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.