skip to content

When the leader of a partition (or queue) dies, what do writers see before another copy is promoted within the same cluster?

level: juniorimportance: must knowfreq 62%

answer

  1. promotion is not instantaneous
  2. detection plus promotion, not promotion alone
  3. clients cache who owns what
  4. rejected is not the same as lost
  5. retry budget must outlast the gap

basics

~20 s

Writes to that partition fail rather than block: the client's cached owner map still names the dead node, so sends are rejected until a replacement copy is promoted. Records survive only if the writer's retries outlast that gap.

solid answer

~40 s

Between the leading node going silent and a replacement being recorded there is a gap in which no node serves writes for that partition (or queue), and that gap is detection time plus promotion time, not promotion alone. Writers do not wait politely: a client caches an owner map — its own view of which node serves which unit of ownership — so it keeps aiming at the dead node and gets connection failures, explicit refusals or timeouts until an error makes it refresh. Once the new leader is recorded and the map is refreshed, sends resume. Nothing is lost as long as the writer's retry budget outlasts the gap; if the caller gives up first, the record was simply never accepted and holding it is the application's problem.

go deeper

for a junior

Remember that a partition has one node serving it, and that losing that node means a short period where sends to it fail. Failed is not the same as lost.

for a middle

Explain the gap as detection time plus promotion time, and explain why the client keeps aiming at the dead node until an error forces it to refresh its cached view of who owns what.

for a senior

Show that you size the writer's retry budget against the worst gap the cluster can produce, and that you measure the gap per partition rather than as a cluster average that hides the slow one.

for a principal

The angle is where you put the cost: a longer client budget hides the gap but lengthens tail latency and grows client buffers, while a shorter one converts routine node loss into user-visible failure.

## The gap is an interval, not an instant A leader-based cluster serves each **unit of ownership** — a partition (or queue) — from exactly one node at a time. When the node holding that leadership dies, the unit has no leader until the cluster records a replacement, and the interval between those two facts is the gap in which no node serves writes for that partition. It has two parts: 1. **Detection.** The silence the cluster tolerates before it treats a node as gone. This is normally the larger part, because it is deliberately set longer than an ordinary pause on a healthy node. 2. **Promotion.** Choosing a copy, recording the new owner, and letting the other nodes and the clients learn of it. On most designs this part is fast, because nothing has to be copied for it to become true. Interviewers ask about this because candidates routinely describe promotion as instantaneous, and then cannot explain why a cluster that survived a node failure still produced a spike of failed sends. ## Why a writer sees errors rather than a pause A client does not ask the cluster where to send every record; it caches an **owner map**, its own view of which node currently serves which unit of ownership, and sends straight there. During the gap that cached map is stale — it still names the node that died — so the writer gets one of: - a connection failure, if the node is simply unreachable; - an explicit refusal, if the node is reachable but no longer serves that unit; - a request timeout, if the connection stands but nothing answers. None of these means the record was lost. They mean it was **not accepted**. A broker has nowhere durable to put a write for a unit it does not lead, so refusing is the honest answer; a cluster that quietly buffered such writes would be promising a durability it cannot deliver. On any of those errors a well-behaved client refreshes its owner map and retries. That refresh is also why the outage the application sees often outlives the promotion: the cluster may record a new leader in a second while the client is still inside a backoff interval, aiming at a node that is gone. ## What decides whether the record survives | What happened to the record | Outcome for the writer | |---|---| | Acknowledged before the leader died, and present on the copy that is promoted | Survives; nothing to do | | Acknowledged, but never reached the copy that is promoted | Can disappear — that is the promotion-time trade, decided by which copy is promoted | | Rejected during the gap, retried once the new leader is known | Accepted normally by the new leader | | Rejected during the gap, and the caller gave up first | Never accepted; the application must hold it or re-derive it | The last row is the one the operator actually controls. The writer's total budget — how long it keeps retrying before it reports failure upward — has to outlast the worst gap the cluster can produce. If that budget is shorter, an ordinary node failure becomes application-visible loss even though the cluster itself lost nothing. There is a second-order effect on the writer's own side. While the gap is open, records accumulate in the client's send buffer. If that buffer fills, the application's send calls begin to block or fail at the source, so a few seconds of gap inside the cluster can surface as a much longer latency excursion in the service producing the records. ## Where designs differ - On platforms that split a stream into parts, each with a leader and a set of copies, the gap is **per unit of ownership**: only the units led by the dead node stop, and the rest of the cluster serves normally throughout. - On platforms whose records live on shared or remote storage — **detached storage** — a replacement needs no copying at all, so nearly the whole gap is detection. - On designs that commit by majority write, any node eligible to take over already holds every committed record, but the unit stays unwritable until a majority is reachable again. - On queue-shaped brokers, where consumers compete for work rather than following a per-part position, the equivalent gap is the time until that queue is served from somewhere else. - Hosted offerings frequently do not expose the detection interval at all, so the cluster half of the gap is the provider's number and the client's budget is the only half you own. ## What to look at afterwards Measure the gap per unit of ownership rather than per cluster: an average over hundreds of units hides the single unit that took minutes. Then check whether writers absorbed the gap or reported failure to their callers — that difference, not the raw duration, is what users actually experienced.

  • Why is the outage the application sees often longer than the time the cluster took to record a new leader?
    Because the client is working from a cached owner map and a backoff schedule. The cluster can record the change in a second while the client is still waiting out an interval before it retries, discovers the error, refreshes its map and reconnects. Reconnection and authentication on the new node add to that, so client behaviour, not the promotion, often dominates the observed outage.
  • Do readers of that partition experience the same gap as writers?
    They lose their source for the same interval, but the consequence differs. A reader that stops for a few seconds simply resumes from where it was once a leader exists, so the gap costs it delay rather than data. On designs that allow reads to be served from a copy that is not the leader, reads may continue through the gap while writes cannot.

saying these in an interview costs you the question

  • Says promotion is instant, so writers never notice anything
  • Believes the cluster buffers writes until a new leader exists
  • Treats a rejected send as a record the cluster lost
  • Assumes the client learns of the new leader without refreshing anything
  • Ignores that the client's own timeout can be shorter than the gap