skip to content

Backpressure at the Broker

What a node does once it cannot keep up: queue the work, delay the answer, block the writer or refuse it, while the client's buffer fills. Asked because the quietest option loses records silently.

part ofBroker & streaming operationsoverview, primer and where to startread it →
on this pageshow

questions

4

A broker node cannot process writes as fast as they arrive, with every configured ceiling respected — what can it do?

level: juniorimportance: must knowfreq 62%

answer

  1. four responses, not one
  2. queue, delay, block, refuse
  3. the call site tells you which
  4. every ceiling honoured, still saturated
  5. the quiet one discards records

basics

~20 s

A saturated node has four responses: queue the work, delay the answer, block the writer, or refuse the request. Each shows up differently at the call site, and the quietest discards records without failing the send.

solid answer

~50 s

Saturation here is not a ceiling being hit — every rate ceiling and connection cap can be honoured by every client and the node still run out of disk, handler threads or network. What it does about that gap falls into four shapes. It can **queue** the work in a bounded internal request queue, which buys time and adds delay. It can **delay** the answer deliberately, so the writer is paced without being told anything is wrong. It can **block** the writer by stopping draining the connection, so pressure lands in the client's own memory and the send call stops returning. Or it can **refuse**, handing the decision back to the application explicitly. Platforms differ over which they pick. The one to watch for is a full queue or buffer configured to discard rather than refuse: the send reports success and the record is gone.

go deeper

for a junior

Recall that a node with more work than it can do has a small set of choices — hold it, answer late, stall the sender, or reject it — and that only rejection is obvious to the caller.

for a middle

Explain why every configured ceiling can be respected and the node still be saturated, and map each of the four responses onto the exact symptom the writing application sees.

for a senior

Show the diagnostic order: name the symptom, identify which response is in play, then reduce offered work before adding capacity. Know that a slow node and a dying node look identical from outside.

for a principal

The question worth owning is which response you want your estate to meet by default, since a silent discard buys availability with data and a refusal buys data with visible failures.

## What "saturated" means here A node is **saturated** when the work offered to it exceeds what it can complete per second. Crucially, this is not the same as a client hitting a configured ceiling. Every client can be inside its `rate ceiling`, every connection cap can be respected and every record can be under the record-size ceiling, and the node can still run out of disk throughput, handler-pool threads or network capacity — because ceilings are configured per client and saturation is the **sum** of all of them plus whatever the node is doing for itself. Backpressure at the broker is the name for whatever the node does about that gap. The reason this is asked at the first screen is that the four available responses produce four completely different symptoms at the call site, and engineers routinely read one as another. ## The four responses - **Queue the work.** The node accepts the request into a bounded internal request queue and processes it when a handler is free. This absorbs a burst well and costs only added delay — until the queue is full, at which point the design must take one of the other exits. - **Delay the answer.** The node serves the request but holds the response longer than it needed to, pacing the writer. Nothing is reported as wrong; the writer simply observes that the call took longer. - **Block the writer.** The node stops draining the connection. Transport-level flow control then pushes the pressure back into the client's own memory: unsent records pile into the **writer send buffer**, and when that is full the send call itself stops returning. - **Refuse.** The node rejects the request with an error the client has to handle. This is the only response that hands the decision back to the application in so many words. ## What each looks like at the call site | Node's response | What the writing application observes | Where the work is waiting | How records can be lost | |---|---|---|---| | Queue | Slower responses, no errors | In the node's request queue | If the queue's overflow exit discards | | Delay | Slower responses, no errors | In the node, deliberately | Not directly; the writer's buffer grows behind it | | Block | Send calls stop returning; threads pile up | In the writer's unsent buffer | If the process dies with the buffer full | | Refuse | An explicit error per request | Nowhere — it was rejected | Only if the application drops what it cannot send | Notice that the first two rows are indistinguishable from each other, and both are indistinguishable from a node that is on its way to failing entirely. That is why "writes got slower" is a starting point for an investigation and never a diagnosis. ## The quiet exit Both the node's request queue and the writer's send buffer are finite, and when either fills the design has a choice: **refuse the caller, or discard the record**. Discarding is the quietest thing in this whole subject. The send call returns successfully. No error is raised, no failure counter on the writing side moves, and a dashboard built on error rates stays green while records disappear. The loss shows up only if somebody reconciles what the application handed to the client against what the cluster actually accepted over the same interval, or reads whatever discard counter the client itself keeps. This is the answer an interviewer is listening for: of the ways to respond to saturation, the one that never fails anything is the one that loses data. ## The response you get varies by platform It is a mistake to learn one platform's habit as the model. Some designs prefer to refuse an over-capacity request outright so that the client retries; some prefer to hold the answer; some rely on transport flow control and never signal anything, so the symptom only ever appears inside the client; a rented cluster may refuse where a self-run one would have queued, because refusing is cheaper to operate at scale. When you describe this in an interview, describe the **behaviour** — queued, delayed, blocked, refused — rather than assuming everyone's node does what yours does. ## What an operator does first 1. **Name the symptom precisely.** "Writes are slower", "writes are being refused", "the sending application has stopped" and "records are missing with no errors" are four different rows of the table above. 2. **Establish which response is in play** before changing anything, because the remedies point in opposite directions: a refusing node wants less offered work or more capacity, while a stalled application wants its own buffer behaviour fixed. 3. **Reduce, then resize, then rewrite.** Take offered work away first (it is the only lever that acts immediately), add capacity second, and change what the writer sends or how it behaves third. One thing not to do reflexively is deepen the node's request queue. A deeper queue converts refusals into latency; the work still has to be done, the node still cannot do it, and now every caller waits longer before finding out.

  • The node answers every request successfully but much later than usual. Has anything been lost?
    Not by that response alone — a delayed or queued request that is eventually answered was still accepted. The exposure is behind it: while answers are slow, records pile into the writer's unsent buffer, and those are lost if the process dies or if the buffer's overflow behaviour is to discard.
  • Why is refusing a request often kinder to the writer than answering it slowly?
    A refusal is explicit, so the application can decide: retry later, slow itself down, divert the record, or fail the user's operation honestly. A slow answer tells it nothing, so it keeps offering the same load, fills its own buffer, and is indistinguishable from a node that is failing outright.
  • Does adding nodes always relieve saturation?
    Only if the work can actually spread. Extra capacity helps when the load is distributed across the cluster, but a single hot stream or one heavy writer can concentrate on the same node, and moving data onto the new node consumes capacity for a while first.

saying these in an interview costs you the question

  • Thinks a saturated node always returns an error to the caller
  • Reads slower writes as a network problem without checking the node
  • Believes queueing the work is free because nothing fails
  • Calls any slowdown a quota even though no ceiling was reached
  • Assumes no exception at the call site means no records were lost
  • Treats a deeper internal request queue as a fix rather than a delay
open as a page

Every thread in a producing application is stuck inside a send call while the broker reports no errors — why?

level: middleimportance: must knowfreq 57%

basics

~20 s

The cluster is accepting records more slowly than the application produces them, so unsent records have filled the writer send buffer. With the buffer full the client blocks the calling thread instead of failing, and broker saturation becomes an application outage.

open as a page

A writer discards records instead of failing the send when its unsent buffer fills under cluster saturation — what has that traded?

level: seniorimportance: should knowfreq 44%

basics

~20 s

Availability of the calling path has been bought with data, and with the evidence that the data is gone. Sends keep succeeding, error rates stay flat, and the loss is detectable only by reconciling records offered against records the cluster accepted.

open as a page

Across many writing applications sharing one cluster, where should overload surface, and what do you require of every writer?

level: principalimportance: should knowfreq 33%

basics

~20 s

Overload should surface in the writers, as bounded waits and explicit failures, rather than in ever-deeper queues on the cluster. Require every writer to bound its unsent buffer, cap how long a send may wait, declare whether its records tolerate loss, and have a plan for records it cannot send.

open as a page