Conflating Price Ticks for Slow Clients Permalink to this section

Part of Market Data & Live Scoreboards, under Real-Time Application Patterns.

Global conflation caps how fast prices leave the feed handler. It does nothing for the viewer on a congested mobile link who cannot absorb even that capped rate. Without per-viewer handling, that viewer’s frames queue in server memory and in socket buffers, and the prices on their screen fall further behind reality every second — the worst possible failure for market data, because the screen looks alive while showing stale numbers. This guide gives every viewer its own conflation: a latest-value buffer that is flushed only when the viewer’s socket can take more.

Symptom & Developer Intent Permalink to this section

  • Some users report prices lagging the market by tens of seconds, growing over the session.
  • Server memory rises with the number of mobile viewers, not with the total number of viewers.
  • Node.js heap snapshots show large arrays of pending write chunks attached to a few response objects.
  • Killing a slow viewer’s connection makes the memory drop immediately.
  • Viewers on good connections are unaffected.

The intent is a stream where every viewer, on any link, sees prices that are at most one flush interval plus network latency behind the market, and where a slow viewer costs the server a bounded amount of memory.

Root Cause Analysis Permalink to this section

When the application writes faster than the network drains, the bytes have to go somewhere. In Node.js, res.write() returns false once the socket’s buffer passes its high-water mark, but it still accepts the data and queues it in user space. Every subsequent write adds to the queue. For a viewer receiving 10 KB per second on a link that drains 4 KB per second, the queue grows by 6 KB every second, forever, and each queued frame is older than the one before.

Price age on a slow viewer's screen Line chart over five minutes of the age of prices displayed to a viewer on a slow link, with queued writes and with per-viewer conflation. Price age on a slow viewer's screen queued writes per-viewer conflation 0 45 90 135 180 0 1 2 3 4 5 minutes connected age of displayed price (s)
With queued writes the lag grows without bound. With per-viewer conflation it stays flat at roughly one drain interval, whatever the market does.

The fix is to stop producing frames for a viewer who has not consumed the previous one. Instead of writing each batch, merge it into the viewer’s own latest-value map and write that map only when the socket has drained. A slow viewer then simply receives fewer frames, each of which is current.

Step-by-Step Resolution Permalink to this section

Step 1 — Give each viewer a latest-value buffer Permalink to this section

// viewer.js — one per open stream.
export class Viewer {
  constructor(res, symbols) {
    this.res = res;
    this.symbols = symbols;
    this.pending = new Map();        // symbol → latest quote not yet written
    this.writing = false;            // true while the socket is backed up
    this.lastId = 0;
  }

  offer(batch) {
    for (const [k, v] of Object.entries(batch.quotes)) {
      if (this.symbols.has(k)) this.pending.set(k, v);   // overwrite: only the latest matters
    }
    this.lastId = batch.id;
    if (!this.writing) this.flush();
  }

  flush() {
    if (!this.pending.size) return;
    const quotes = Object.fromEntries(this.pending);
    this.pending.clear();
    const ok = this.res.write(`event: batch\nid: ${this.lastId}\ndata: ${JSON.stringify({ quotes })}\n\n`);
    if (!ok) {
      this.writing = true;
      this.res.once('drain', () => { this.writing = false; this.flush(); });   // catch up in one frame
    }
  }
}

The memory held per viewer is bounded by the number of symbols it watches, not by time or market activity. A viewer that stalls for a minute receives one frame with the latest value of each symbol when it recovers.

Per-viewer conflation breaks batch-to-batch continuity by design, so a prev field would report a gap on every merged frame. Rely instead on per-key source sequences: the client applies a quote only if its sequence is newer than the one it holds. Continuity at the batch level matters for replay, and a state-shaped feed does not replay.

es.addEventListener('batch', (e) => {
  const { quotes } = JSON.parse(e.data);
  for (const [k, v] of Object.entries(quotes)) {
    const cur = board.get(k);
    if (!cur || v.s > cur.s) { board.set(k, v); changed.add(k); }
  }
});

Step 3 — Wire viewers to the batch feed Permalink to this section

const viewers = new Set();
edgeBus.on('batch', (b) => { for (const v of viewers) v.offer(b); });

app.get('/api/quotes/stream', (req, res) => {
  const symbols = parseSymbols(req.query.s);
  openStream(res, { retryMs: 1000 });
  const v = new Viewer(res, symbols);
  res.write(`event: snapshot\ndata: ${JSON.stringify({ quotes: pick(edgeState.current(), symbols) })}\n\n`);
  viewers.add(v);
  req.on('close', () => viewers.delete(v));
});
A slow viewer receives fewer, fresher frames Sequence diagram showing three batches arriving while a viewer's socket is backed up, being merged into its pending map, and a single merged frame written when the socket drains. A slow viewer receives fewer, fresher frames Batch feed Viewer buffer Socket batch 501 ACME 101.20 write → false (backed up) batch 502 ACME 101.22 batch 503 ACME 101.19, GLOBX 44.12 drain one frame: ACME 101.19, GLOBX 44.12
Three upstream batches became one downstream frame, and every value in it was current when it was written.

Step 4 — Protect against viewers that never drain Permalink to this section

A connection that stays backed up for a long time is probably dead. Close it after a deadline and let EventSource reconnect; the snapshot on reconnect repairs everything.

const STALL_MS = 30_000;
offer(batch) {
  // …merge as before…
  if (this.writing && !this.stalledSince) this.stalledSince = Date.now();
  if (!this.writing) this.stalledSince = null;
  if (this.stalledSince && Date.now() - this.stalledSince > STALL_MS) this.res.destroy();
}

Step 5 — Tell the viewer when it is behind Permalink to this section

Freshness should be visible. Include the batch’s upstream timestamp and let the client show a “delayed” badge when the newest data is older than a threshold:

// Server: add ts (upstream receive time) to each written frame.
// Client:
const age = Date.now() - frame.ts - clockOffset;
delayedBadge.hidden = age < 2000;

For regulated market data this is not only good practice: displays often must mark delayed prices explicitly.

Validation & Monitoring Permalink to this section

# Simulate a slow link and watch memory and frame count on the server.
curl -sN --limit-rate 3k 'https://md.example.com/api/quotes/stream?s=ACME,GLOBX,INIT,ZETA' > /dev/null &
watch -n 5 'curl -s localhost:9100/metrics | grep -E "sse_viewer_pending_keys|process_resident_memory"'

With per-viewer conflation, sse_viewer_pending_keys for the slow viewer never exceeds the number of symbols it watches, and resident memory stays flat.

Server memory after 10 minutes with 500 slow mobile viewers Bar chart of resident memory for queued writes versus per-viewer conflation with 500 viewers on slow links. Server memory after 10 minutes with 500 slow mobile viewers Queued writes 3.1 GB and rising Per-viewer conflation 180 MB, flat resident memory of the edge process
Queued writes store the market's history for each slow viewer. Conflation stores one value per symbol.

Export per-viewer metrics sparingly — aggregate them. Useful ones are the distribution of frames written per viewer per second (a slow cohort shows as a low tail), the number of viewers currently backed up, and stall-closes per minute.

Production Checklist Permalink to this section

Frequently Asked Questions Permalink to this section

Is per-viewer conflation expensive with many viewers?

It costs one small map per viewer and a merge per batch. That is far cheaper than queueing frames, and fast viewers — the majority — never accumulate anything because their socket is never backed up.

Can the same approach work in Go or Python?

Yes. In Go, give each viewer a goroutine with a mutex-protected map and a notify channel of capacity one; in Python asyncio, a dict plus an asyncio.Event that the writer task waits on. The rule is the same: merge while the writer is busy, write the merged state when it is free.

Does this lose data?

It discards intermediate values, which for prices and scores carry no information once superseded. Discrete events such as trades must never go through this path; send them as a separate, replayable event type.

How do I detect slow viewers before they cause trouble?

Watch res.writableNeedDrain or the socket's buffered size. A viewer that is backed up on most batches is slow; with conflation that is harmless, and the metric simply tells you how many users are on constrained networks.