skip to content

A writer's cached owner map still names the superseded leader for a partition (or queue): why does its next write not quietly succeed?

level: seniorimportance: should knowfreq 47%

answer

  1. the cached route goes stale routinely
  2. the danger is silent success
  3. authority checked where durability is decided
  4. refusal is a routing signal
  5. refresh the map, then retry

basics

~20 s

Because authority is checked where the write must become durable, not where the client sends it. The stale route is refused under an out-of-date generation number, the error travels back, and the client's cached ownership view is corrected rather than the write silently landing on a node nobody reads from.

solid answer

~50 s

A client keeps an **owner map** — its cached view of which node serves which partition (or queue) — and that view goes out of date the moment leadership moves. Sending to the old node is therefore normal, not exceptional. What must not happen is a silent success: a write that the old node accepts locally, answers cheerfully, and that no reader ever sees. Two things prevent it. First, the node itself may already know it was replaced and refuse. Second, and this is the case that matters, even a node that still believes it leads cannot get the write accepted by the copies or storage that would make it durable, because it carries an older generation number. The refusal surfaces to the client as a not-the-owner error, and the client's correct response is to refresh its owner map and retry against the current owner — not to retry the same node harder.

go deeper

for a junior

Take away one fact: a client caches which node owns a partition (or queue), that cache goes out of date, and the broker refuses the misrouted write instead of quietly accepting it.

for a middle

Explain why the refusal does not depend on the old node knowing anything: it carries an out-of-date generation number, so it cannot get the write accepted where durability is decided.

for a senior

Demonstrate production judgment — read the error burst after a leadership change as normal, investigate one that does not subside, spot a writer pinned to a single address, and treat a timeout as unknown rather than failed.

for a principal

Set the expectation for the estate: every writing service must re-resolve the owner as part of its retry, and any path into the data that skips the ownership check needs a named owner and a reason to exist.

## The situation A client that writes to a broker does not look up the owner of a partition (or queue) on every call; it caches an **owner map**, the view of which node currently serves which unit of ownership. Leadership can move at any moment — a node failed, a node was restarted, a leadership-only rebalance ran — and no design pushes that change into every client instantly. So a client sending to a node that is no longer the owner is an ordinary, expected event, happening quietly many times a day in a healthy cluster. The danger is not the stale route. It is the stale route **succeeding**. ## The outcome that must not happen Imagine the old node accepts the write, appends it to its local storage, and answers the client that all is well. The client moves on. Readers, meanwhile, are being served by the current owner, which never saw that record. Nothing errors anywhere. You would discover it as records that were "definitely sent" and are simply not in the stream, which is one of the most expensive classes of incident to diagnose because the writing side has evidence of success. ## Why it does not happen | Where the write goes | What happens | What the client sees | |---|---|---| | To the current owner | accepted and acknowledged normally | success | | To a node that already learned it was replaced | refused immediately | a not-the-owner error | | To a node that still believes it leads | it cannot get the write accepted where durability is decided, because it carries an older generation number | an error, after a delay | | To a node that is unreachable from the client | no answer | a timeout, which is an unknown outcome, not a failure | The third row is the important one, because it does not depend on the old node being well-informed. Authority is checked **where the write has to be accepted by someone else** — the copies that must take it, or the shared storage that holds it. A node with an out-of-date generation number can believe whatever it likes; it cannot make a record durable, and it cannot manufacture an acknowledgement that means anything. ## The correction path 1. The client's write is refused, or its acknowledgement never comes. 2. The client treats a not-the-owner refusal as a routing signal rather than a transport fault, and refreshes its owner map. 3. It retries against the node the corrected map names. 4. If that one is also stale — leadership can move twice during an incident — the cycle repeats, which is why a bounded retry with backoff matters rather than an immediate tight loop. Designs differ in how step 2 is triggered: some brokers return an error that names the condition, some redirect the client to the current owner, and some hand back a refreshed view with the error. What is common across them is that the client must react to it, and that the correction is **pulled on failure** rather than pushed on change. ## What this asks of the client side - A retry that goes to **the same address** forever turns a routing problem into an outage. The retry must include re-resolving the owner. - A client pinned to a single node's address rather than to the cluster loses the correction path entirely; this is a real configuration mistake, not a theoretical one. - A **timeout** is not a failure. If the acknowledgement never arrived, the record may or may not exist, and whatever the application does next has to be safe under both. - "The library handles it" is usually true and worth confirming rather than assuming, especially for a writer someone built directly against the wire. ## What it asks of the operator - Expect a burst of not-the-owner errors after any leadership change. They are the mechanism working, and alerting on them as failures teaches everyone to ignore a useful signal. What deserves attention is a burst that **does not subside**, because that means clients are not correcting their view. - During a change, watch for one writer stuck on an old route while others recovered — it usually points at that client's retry behaviour or its configuration rather than at the cluster. - Be suspicious of any path into the data that bypasses the ownership check — a repair tool, an import, a restore — because nothing on that path can be refused for carrying stale authority. ## Where designs differ On platforms where a stream is split into parts and each part has one owner, the client's owner map is per part and the correction is per part too. On queue-shaped brokers where consumers and producers attach to whichever node fronts the queue, the same staleness exists but is often hidden behind a connection-level redirect. On designs with detached storage, ownership can move without copying anything, so the map goes stale **more** often, not less — the cheapness of the move is exactly what makes it frequent.

  • Should an operator alert on not-the-owner errors?
    Not on their presence — a burst follows every leadership change and is the mechanism working. Alert on a burst that does not subside, or on one writer still producing them after others recovered, which points at that client's retry or configuration rather than at the cluster.
  • What is the worst thing a client can do when it gets that refusal?
    Retry the same address in a tight loop. It converts a routing correction into a self-inflicted outage and hides the real condition. The retry must re-resolve which node owns the unit, and back off between attempts so a cluster-wide change does not become a stampede.
  • Why can leadership moving more often make stale routes more common rather than less?
    Because where a unit's records live on shared or remote storage, moving ownership copies nothing and is therefore cheap enough to do routinely. Cheap moves happen often, and every move leaves every cached owner map briefly wrong — which is fine as long as clients correct on refusal.

saying these in an interview costs you the question

  • Assumes a client always learns a leadership change before it writes
  • Says a node's local acceptance means the write is safe
  • Believes the client's cached ownership view can never be out of date
  • Treats a not-the-owner refusal as a transport fault and retries the same node
  • Reads a timeout as a definite failure rather than an unknown outcome