skip to content

A broadcast hub pushes each room message to thousands of WebSocket subscribers, and one slow subscriber delays everyone - where should the buffer sit, and why?

level: seniorimportance: should knowfreq 42%

answer

  1. the slowest reader sets the pace
  2. one waiting place couples everyone
  3. a queue per connection, not per room
  4. encode once, enqueue a reference
  5. overflow costs only the laggard

basics

~20 s

Per subscriber, not shared: encode each message once, put a reference on each subscriber's own bounded queue, and drain each queue with its own writer. Overflow then drops, keeps only the latest, or disconnects only the laggard.

solid answer

~50 s

If the hub writes to each subscriber in turn and waits for each write, or drains one shared queue that way, everything runs at the pace of the slowest connection - and an overflow policy on that shared queue costs **every** subscriber messages. So the buffer moves to the subscriber: the hub serialises each message once and offers a reference to a **bounded send queue per subscriber**, and a separate writer drains each queue at its own connection's speed. When one queue fills, the overflow policy fires for that subscriber only - drop new messages, keep only the latest for state-like feeds, or disconnect the slow consumer so it reconnects and resynchronises. I would watch per-subscriber queue depth percentiles, overflow and disconnect counters, and the time to fan one message out, which should stay flat when one client slows.

code

pseudocode · 22 lines
pseudocode
function broadcast(room, message):
    frame = encode(message)                  // serialised once, shared
    for sub in room.subscribers:
        offer(sub, frame)                    // never waits on a socket

function offer(sub, frame):
    if size(sub.queue) < sub.capacity:
        append(sub.queue, frame)             // a reference, not a copy
        return
    recordOverflow(sub)
    if sub.policy == KEEP_LATEST:
        clear(sub.queue)
        append(sub.queue, frame)             // only the newest state survives
    else if sub.policy == DROP_NEWEST:
        discard(frame)                       // this subscriber misses it, nobody else
    else:                                    // DISCONNECT
        closeWithDeadline(sub.connection)    // client reconnects and resyncs

function writer(sub):                        // one per connection
    while isOpen(sub.connection):
        frame = takeOldest(sub.queue)        // waits while the queue is empty
        write(sub.connection, frame)         // may wait; holds up only this subscriber

go deeper

for a junior

Recall the core picture: one slow receiver must not hold up the others, so each subscriber gets its own small queue, and whatever overflows is that subscriber's loss alone.

for a middle

Explain why a loop that waits on each write, or one shared queue drained that way, couples every subscriber, and how encoding once and enqueuing a reference per bounded queue decouples them.

for a senior

Show the operating judgment: pick drop, keep-latest or disconnect per feed, bound queues in bytes, close connections whose closing frame cannot drain, and alert on depth percentiles and slow-consumer disconnects.

for a principal

Weigh the costs: per-subscriber capacity times connections per node is a memory budget, disconnecting turns slow clients into reconnect load, and each feed should state what it may lose.

## How one slow subscriber stalls everyone A **broadcast hub** is the component that takes one message for a room, channel or topic and pushes it to every live subscriber over a long-lived connection - a WebSocket or a Server-Sent Events stream. Subscribers drain at very different speeds: most sit on good links, a few on a congested mobile network or a frozen browser tab. Two common designs tie every subscriber to the slowest one: - **A synchronous send loop.** The hub walks the subscriber list and writes the message to each connection, waiting until each write is accepted. When one subscriber stops reading, its connection's send buffer fills, the transport's own flow control holds the write, and the loop stops on that subscriber. Every subscriber after it, and every later message, waits. - **One shared queue in front of the loop.** A single queue between publisher and loop only moves the stall. The loop still drains that queue at the pace of the slowest write, so the queue fills, and the overflow policy now fires on the **shared** queue: dropping or conflating there costs every subscriber messages because one of them is slow. Separate connections do not help on their own. What couples the subscribers is the single place where every delivery waits. ## Where the buffer belongs The fix is to put the buffer, and with it the overflow decision, **per subscriber**: 1. **Encode once.** Serialise the message into its wire frame a single time per broadcast. The bytes are immutable and shared. 2. **Enqueue a reference per subscriber.** The hub offers the frame to each subscriber's own **bounded send queue**. Offering never waits on a socket; it either succeeds or hits that subscriber's overflow rule. 3. **Drain each queue at its own pace.** A writer per connection (a task, a readiness-driven write, whatever the runtime offers) takes the oldest frame and writes it. If that write waits, only that subscriber waits. Fan-out becomes one encode plus one constant-time offer per subscriber, and its duration no longer depends on how fast any single client reads. ## Choosing the policy per subscriber Once each subscriber has its own queue, the overflow policy is applied to that queue alone. Which one depends on what a single message in the feed means: | policy at a full queue | fits | what that subscriber loses | |---|---|---| | **Drop the newest** | low-value notifications | the messages that arrive while it is full, silently | | **Keep only the latest** | state that each update supersedes: a score, a position, a price | intermediate values a state feed does not need | | **Disconnect the slow consumer** | feeds where every message matters, such as chat | its connection; the client reconnects and resynchronises | Whichever applies, the loss is **contained** to the subscriber that could not keep up; everyone else keeps receiving on time. Disconnecting is often the honest choice for feeds where gaps are unacceptable: a silent gap is worse than a visible reconnect, and the client can catch up from a snapshot or history (how it does that belongs to delivery and replay, not to the overflow policy). One practical detail: on a WebSocket, a graceful closing frame queues behind the same backlog the client is not reading, so hubs typically send it with a short deadline and then close the underlying connection outright. ## What it costs and where the bound goes Isolation is not free; it converts a stall into a memory budget: - **Bound every queue.** An unbounded per-subscriber queue isolates latency but not memory: one stuck client grows its queue until the process runs out. - **Bound in bytes as well as count.** Message sizes vary, and a count bound alone lets a few large frames blow the budget. - **Budget the worst case.** Capacity per subscriber times subscribers per node is the most the node can hold queued. - **Remember what a reference keeps alive.** Sharing the frame saves copies, but a lagging subscriber's queue keeps old frames in memory until it drains or is cut off. ## What to measure The failure is per subscriber, so the signals must be too: - **Per-subscriber queue depth** as a distribution - maximum and high percentiles, never only the average, which hides the one client falling behind. - **Overflow counters per subscriber and per room**: drops, conflations and slow-consumer disconnects, tagged by reason. - **Time in queue**, from enqueue to write, which shows how stale delivered messages are. - **Fan-out duration**, the time to offer one message to every subscriber. It should stay flat; if it rises when a single client slows, isolation is broken and something is still waiting on a socket. ## How to answer it out loud 1. Name the coupling: a loop that waits on each write, or one shared queue, runs at the slowest subscriber's pace. 2. Move the buffer: encode once, a bounded queue per subscriber holding references, a writer per connection. 3. Apply the overflow policy per subscriber, chosen by what the data can afford to lose. 4. Say how you would see it: per-subscriber depth percentiles, overflow and disconnect counters, flat fan-out time. The idea the interviewer listens for: the buffer belongs where the slowness lives, so a slow client's cost lands on that client.

  • Should a chat room and a live score feed use the same per-subscriber overflow policy?
    No. Each score update supersedes the last, so keeping only the latest loses nothing the client wanted. Every chat message matters, so silently dropping one leaves an invisible gap; disconnecting the slow client is the better choice, because it reconnects knowing it must catch up from history. The policy follows what one message means, which is why it is set per feed, not per hub.
  • How do you choose the capacity of each subscriber's queue?
    Size it for the longest stall a healthy client should survive - a few seconds of mobile hiccup at the feed's peak rate - and bound it in bytes as well as count, because frame sizes vary. Then multiply by the subscribers a node holds to check the worst case fits in memory. Bigger is not safer: it only delays the overflow and makes queued messages staler.
  • Why not simply send the slow subscriber a closing frame and forget it?
    On a WebSocket the closing frame is queued behind the same unread backlog, so a client that is not reading may never see it, and the connection lingers holding buffers. Send it with a short deadline, then close the underlying connection outright and release that subscriber's queue.

saying these in an interview costs you the question

  • Thinks separate connections alone stop one slow subscriber delaying the rest
  • Applies the overflow policy to one shared queue, so one laggard costs everyone messages
  • Blocks the publisher until the slowest subscriber's queue has room
  • Calls an unbounded per-subscriber queue isolation, ignoring that one stuck client exhausts memory
  • Watches average queue depth, which hides the single subscriber falling behind