A WebSocket server that pushes data to many clients will eventually meet one that cannot keep up: a phone on a weak cell link, a laptop that went to sleep with the tab open, a dashboard whose message handler takes 50 ms per update. The server can produce faster than that client can consume, and the difference has to go somewhere. If it goes into an unbounded buffer the process runs out of memory; if it blocks the publisher, every other client stalls too.

This page is about that slow consumer on the server side of a fan-out: where the bytes pile up, how each runtime tells you, which policy to apply when a client falls behind, how to spot a stalled reader early and how the client should recover. Credit-based flow control for one-to-one streams and relays is covered in backpressure in bidirectional streams.

Advertisement

WebSocket has no flow control of its own

RFC 6455 defines frames, masking, fragmentation and control messages, but no window and no credit mechanism. The only flow control on a classic WebSocket connection is the one underneath it: TCP. When the receiving application stops reading, its kernel receive buffer fills, the TCP window it advertises shrinks to zero, and the sender's kernel stops transmitting. The sender's kernel send buffer then fills, and only after that does the sending application notice anything at all.

That chain works, but it reacts late: kernel buffers on both ends can hold hundreds of kilobytes or more before your code sees a signal, and the signal is a write that waits or a buffered-byte counter that grows.

SERVERPublisherone eventPer-client queueyour policy lives hereLibrary buffere.g. write_limitKernel send bufferSO_SNDBUFTCP windowNETWORKKernel receive bufferadvertises windowLibrary inbound queuee.g. max_queueApp handleronmessage / recvCLIENTWhen the handler is slow, every buffer to its right fills in turn, then the TCP window closes,the server kernel buffer fills, the library buffer passes its high-water mark, and send() starts to wait.Only the per-client queue is under your control. Bound it in bytes and decide what happens when it is full.
The buffers between a publisher and a slow client handler. TCP backpressure fills them right to left; the per-client queue is the only one whose policy you choose.

How each runtime surfaces the signal

RuntimeSend-side signalReceive side
Browser WebSocketbufferedAmount: bytes queued by send() but not yet handed to the networkNo backpressure: onmessage fires for every message whether you are ready or not
Browser WebSocketStreamStreams: await writer.write() waitsReading from a ReadableStream pulls on demand. Chrome only at the time of writing; feature-detect it
Python websockets 16await ws.send() waits while the transport buffer is above write_limit (default 32 KiB high-water mark)max_queue (default 16 frames) is the high-water mark for unread incoming frames; above it the library stops reading the socket until the queue drains
Node wsws.bufferedAmount and the send() callbackMessages are emitted as they arrive; pause the underlying socket yourself if needed
Go gorilla/websocketWriteMessage blocks; SetWriteDeadline bounds itYou call ReadMessage in your own loop, so reading is naturally pull-based

Two patterns stand out. Pull-based APIs (Go, Python, WebSocketStream) give backpressure for free on the receive side, because nothing is read until you ask. Push-based APIs (the browser WebSocket, Node event emitters) hand you messages regardless, so a slow handler builds its own invisible queue in memory.

Advertisement

The policies

When a client's queue is full there are five things a server can do, and only one is never acceptable in a fan-out.

  • Block the publisher. Correct for a one-to-one stream, where slowing the producer is the point. Fatal in a fan-out: one slow client stalls delivery to every other client sharing that publisher loop. This is head-of-line blocking, and it is the most common slow-consumer bug.
  • Buffer without limit. Postpones the problem until the process runs out of memory.
  • Drop. Discard the newest or the oldest message. Fine for data where any gap is tolerable and the next message supersedes nothing in particular, such as a chat typing indicator.
  • Conflate. Keep only the latest value per key. Ideal for state feeds such as prices, positions and presence, where an old update for a key is worthless once a newer one exists. The queue is bounded by the number of distinct keys, not by time.
  • Disconnect. Close the connection and let the client resynchronise. The right default for feeds where every message matters, because a client that is permanently behind is better served by a fresh snapshot than by a backlog it will never clear.

Worked example: three clients, four policies

The simulation below publishes one update per tick for four instruments to three clients. The fast client drains two messages per tick, the slow one half a message, and the stalled one nothing. Each client has its own queue, limited to eight messages where the policy uses a limit.

from collections import deque

def run(policy, ticks=60, limit=8):
    """A hub publishes 1 message per tick to three clients that drain at
    different rates. Each client has its own outbound queue."""
    rates = {"fast": 2.0, "slow": 0.5, "stalled": 0.0}
    queues = {name: deque() for name in rates}
    credit = {name: 0.0 for name in rates}
    dropped = {name: 0 for name in rates}
    closed = {}
    for t in range(ticks):
        for name, q in queues.items():
            if name in closed:
                continue
            msg = (f"sym{t % 4}", t)        # four instruments, one update per tick
            if policy == "conflate":
                # Keep only the latest value per key: replace a queued update for the same key.
                for i, (key, _) in enumerate(q):
                    if key == msg[0]:
                        del q[i]; dropped[name] += 1
                        break
                q.append(msg)
            elif policy == "unbounded" or len(q) < limit:
                q.append(msg)
            elif policy == "drop-oldest":
                q.popleft(); q.append(msg); dropped[name] += 1
            elif policy == "disconnect":
                closed[name] = t; q.clear()
                continue
            credit[name] += rates[name]
            while credit[name] >= 1 and q:
                q.popleft(); credit[name] -= 1
            credit[name] = min(credit[name], 1.0)
    for name in rates:
        state = f"closed at tick {closed[name]}" if name in closed else f"queue={len(queues[name]):2}"
        print(f"  {policy:11} {name:8} {state:20} dropped={dropped[name]}")

for policy in ("unbounded", "drop-oldest", "conflate", "disconnect"):
    run(policy)

Running it prints:

  unbounded   fast     queue= 0             dropped=0
  unbounded   slow     queue=30             dropped=0
  unbounded   stalled  queue=60             dropped=0
  drop-oldest fast     queue= 0             dropped=0
  drop-oldest slow     queue= 7             dropped=23
  drop-oldest stalled  queue= 8             dropped=52
  conflate    fast     queue= 0             dropped=0
  conflate    slow     queue= 3             dropped=27
  conflate    stalled  queue= 4             dropped=56
  disconnect  fast     queue= 0             dropped=0
  disconnect  slow     closed at tick 15    dropped=0
  disconnect  stalled  closed at tick 8     dropped=0

With no limit, the queues grow linearly until memory runs out. Drop-oldest caps memory at the limit but the slow client loses 23 of 60 updates with no idea which. Conflation keeps the queue at the number of keys, four, and every message it drops was already superseded, so the slow client still ends with the latest value of every instrument. Disconnect protects the server soonest and removes the stalled client at tick 8, the moment its queue overflows. In every policy the fast client is untouched, because no client's queue can delay another's.

A bounded per-client queue in Python

The server below uses the Python websockets library. The publisher never awaits anything per client: it serialises once, offers the frame to every client's bounded queue with put_nowait, and closes any client whose queue is full. Each client has its own writer task, and that task is the only place that waits on the network.

import asyncio
import json

from websockets.asyncio.server import serve

QUEUE_LIMIT = 256          # messages per client before we act
SEND_TIMEOUT = 10.0        # seconds one send may wait on a full transport
SLOW_CONSUMER = 4001       # private-use close code (RFC 6455 reserves 4000-4999)

clients = set()


class Client:
    def __init__(self, ws):
        self.ws = ws
        self.queue = asyncio.Queue(maxsize=QUEUE_LIMIT)

    def offer(self, frame):
        """Called by the publisher. Never blocks the publisher."""
        try:
            self.queue.put_nowait(frame)
            return True
        except asyncio.QueueFull:
            return False

    async def writer(self):
        while True:
            frame = await self.queue.get()
            # send() waits while the transport buffer is above write_limit,
            # so this is where TCP backpressure reaches our code.
            send = asyncio.ensure_future(self.ws.send(frame))
            send.add_done_callback(lambda t: t.cancelled() or t.exception())  # always read
            done, _ = await asyncio.wait({send}, timeout=SEND_TIMEOUT)
            if not done:
                # The websockets docs discourage cancelling send(); close instead.
                await self.ws.close(SLOW_CONSUMER, "send timeout")
                return
            if send.exception():   # connection closed: stop this writer
                return


async def handler(ws):
    client = Client(ws)
    clients.add(client)
    task = asyncio.create_task(client.writer())
    try:
        await ws.wait_closed()
    finally:
        clients.discard(client)
        task.cancel()


def publish(event):
    frame = json.dumps(event)          # serialize once, not once per client
    for client in list(clients):
        if not client.offer(frame):
            clients.discard(client)
            asyncio.create_task(client.ws.close(SLOW_CONSUMER, "slow consumer"))


async def main():
    async with serve(handler, "0.0.0.0", 8765, write_limit=64 * 1024):
        await asyncio.Future()   # publish() is called from your event source

Three details matter. The close code 4001 is in the 4000-4999 range RFC 6455 reserves for private use, so a client that is merely slow can tell a slow-consumer disconnect from a normal close. A client stalled at the TCP level never sees that code: the close frame sits behind the backlog, the library gives up after close_timeout and aborts the connection, and the client observes 1006. Clients should treat both as a cue to resync. Second, the writer does not cancel a pending send(), which the websockets documentation discourages; when a send has waited ten seconds it closes the connection instead. Third, the limit counts messages, which is only safe if messages are small and similar; count bytes if sizes vary.

The Go chat example in gorilla/websocket uses the same shape: each client has a buffered send channel of 256 messages, the hub sends with a non-blocking select, and on the default branch it closes that client's channel and removes it. Its write pump sets a 10-second write deadline. If you write Go, start there.

Detecting a stalled reader early

Queue depth is a lagging signal. Better ones, in rough order of usefulness:

  • Age of the oldest queued message. It tells you directly how far behind real time the client is, independent of message rate. Disconnect when it passes your freshness budget, such as five seconds for a trading screen.
  • Time blocked in send. A write that has waited for seconds means the TCP window is closed: the client is not reading, or its network has stopped.
  • Missing pongs. A client that does not read cannot answer pings. Python websockets defaults to ping_interval=20 and ping_timeout=20 seconds. Keepalive mostly catches dead connections, not slow ones, but it is free; ping/pong keepalive covers tuning it.
  • Library buffer growth. bufferedAmount or the transport's buffered byte count rising for several seconds.

When the browser is the slow side

Often the network is fine and the client is the bottleneck: the message handler parses large JSON, updates a framework store, and triggers a render for every message. Because the browser WebSocket pushes every message, the backlog builds in the page's memory and the tab gets slower as it falls behind.

// Browser: you can see your own send backlog, not the server's.
const ws = new WebSocket("wss://example.com/feed");
const HIGH_WATER = 1 << 20;   // 1 MiB queued in the browser

function trySend(msg) {
  if (ws.bufferedAmount > HIGH_WATER) return false;   // caller retries or drops
  ws.send(msg);
  return true;
}

// Receive side: keep onmessage cheap. Store the latest value, render once per frame.
const latest = new Map();
ws.onmessage = (e) => { const m = JSON.parse(e.data); latest.set(m.key, m); };
function render() { for (const m of latest.values()) draw(m); latest.clear(); requestAnimationFrame(render); }
requestAnimationFrame(render);

The receive pattern is client-side conflation: store the latest value per key in onmessage, and render once per animation frame. If parsing itself is the cost, move the socket into a Web Worker and post only the conflated results to the main thread. For uploads, check bufferedAmount before each send(), because sending without checking queues unbounded data in the browser.

Recovering after a slow-consumer disconnect

Disconnecting only helps if the reconnect does not recreate the backlog. If the client reconnects and asks for everything since its last sequence number, a client that could not keep up before now has a larger backlog than ever. Instead, on a reconnect after code 4001, send a snapshot of current state, then resume the live stream from the snapshot's sequence number. Clients should reconnect with exponential backoff and jitter, as described in WebSocket reconnection strategies, so that a server-side hiccup that marks many clients slow does not bring them all back at once.

Failure modes

  • Awaiting send inside the broadcast loop. One slow client delays every client. Look for for client in clients: await client.send(...) in code review.
  • Serialising per client. CPU cost multiplies by the connection count; serialise once and share the bytes.
  • Compression multiplying memory. permessage-deflate keeps per-connection compression state, which adds memory to every slow client you are holding on to; see permessage-deflate.
  • Limits in messages when sizes vary. A limit of 256 messages is 256 KB for small updates and hundreds of megabytes for snapshots.
  • Silent drops. Dropping without telling the client leaves it showing wrong data. Add sequence numbers so the client can detect gaps and resync.

Sizing it

Do the worst-case arithmetic before launch. With 50,000 connections, a 256-message queue and 1 KB messages, a server holding every queue full needs about 12.8 GB, before library and kernel buffers. That is why per-client limits must be small and why disconnecting is often the right answer: memory spent on a client that will never catch up is memory taken from the clients that can. Export the distribution of queue depth and oldest-message age per server, and count disconnects by reason so a spike in slow-consumer closes is visible.

What to do next

  1. Find every broadcast loop in your server and confirm it never awaits a per-client send.
  2. Give every connection a bounded outbound queue, measured in bytes, with an explicit policy: drop, conflate or disconnect.
  3. Choose conflation keys for state feeds and add sequence numbers to everything else.
  4. Add a send timeout and an oldest-message-age check, and disconnect with a private-use close code.
  5. Make reconnects after a slow-consumer close start from a snapshot, with jittered backoff.
  6. In browser clients, keep onmessage cheap and render at most once per frame.
  7. Load-test with deliberately stalled clients and confirm fast clients see no added latency.
Key takeaway: WebSocket inherits its only flow control from TCP, so a slow reader shows up late, as a send that waits or a buffer that grows. In a fan-out, never let that wait reach the publisher: give each client its own bounded queue and a writer task, serialise once, and decide per feed whether a full queue drops, conflates by key or disconnects. Watch message age and send time, not only depth, and make reconnects start from a snapshot so a disconnect actually clears the backlog.