skip to content

Why does leadership pile onto a few cluster nodes after failures and restarts, and why does it stay there?

level: middleimportance: should knowfreq 52%

answer

  1. reassigned on loss, not reclaimed on return
  2. who was alive at the time
  3. the plan's first-listed node is home
  4. restoration is an explicit act

basics

~20 s

Every stop hands a unit's leadership to a surviving copy that is caught up, and a node that comes back rejoins as an ordinary copy. Leadership accumulates on whichever nodes happened to be alive at the wrong moments and stays until something hands it back.

solid answer

~50 s

Leadership moves on failure and does not walk home on its own. When a node stops, every unit of ownership it was serving — a partition, or a queue on queue-shaped brokers — is handed to another copy that is current enough to serve. When that node returns it rejoins as an ordinary copy and catches up; nothing in that sequence gives it back what it was serving. Repeat this over a few node failures and one rolling restart and the assignment is shaped by the accident of who was alive each time. Each unit has a **home node** — the node listed first in its plan, where leadership is expected to sit — and drift is measured against that. Designs differ in whether restoration runs on its own, on a schedule, or only when an operator asks for it.

go deeper

for a junior

Recall the asymmetry: when a node stops, another copy has to start serving its units immediately; when it comes back, nothing forces the work to return. That one-way movement is the whole cause.

for a middle

Explain the home node as the first node in a unit's plan, distinguish being caught up from being the serving node, and name the events that concentrate leadership: repeated failure of one node, a fixed restart order, a batch stopped together.

for a senior

Demonstrate that you treat restoration as an operation with a time and a cost rather than as something the cluster does for you, and that you check the per-node view after every planned change rather than trusting that the cluster settled.

for a principal

Consider the estate: if drift is invisible by default and correction is manual, skew is a standing tax on headroom. Decide whether restoration is automatic with guardrails, and what evidence a change must produce before it is called finished.

## What moves when a node stops On designs where exactly one node serves each **unit of ownership** — a partition, or on queue-shaped brokers a queue — stopping that node forces a choice for every unit it was serving. Another copy that is current enough takes over, the change is recorded, and clients refresh their **owner map** (their cached view of which node serves what) on the next miss. No records move: the copies were already there, which is exactly why the handover is fast enough to be survivable. The selection is constrained by who is available and current at that instant, not by any plan for even distribution. That is the seed of skew. ## The home node Each unit has a plan naming the nodes that should hold a copy of it. The first node in that list is its **home node**: the one leadership is expected to sit on when nothing is wrong. The home nodes are chosen so that, taken across the whole cluster, leadership would be spread evenly. That expectation is only realised while nothing has failed. The plan describes where leadership *belongs*; the live assignment describes where it *is*. Skew is the gap between those two, and the gap only ever opens in one direction on its own. ## Why it does not come back by itself When a stopped node returns, three things happen: its process starts, its copies fetch whatever they missed until their **catch-up distance** reaches zero, and they rejoin the **caught-up set** — the copies current enough to be promoted or to satisfy a write. None of those steps includes taking back the units this node used to serve. There is no penalty for leaving leadership where it is; the units are being served correctly, just not from home. So the return to home must be an explicit act, and designs differ in who performs it: - some run a **restoration continuously or periodically** on their own, correcting drift without anyone asking; - some expose it as an **operator action** that has to be run deliberately; - some never concentrate this way at all, because no single node serves a unit — where consumers compete for work from a shared queue there is nothing to hand back. Saying which of those you mean is the difference between a correct answer and one platform's answer. ## The patterns that concentrate it 1. **A node that fails repeatedly.** Each of its failures sheds leadership onto its peers; each recovery gives none back. A flapping node ends up serving almost nothing while its peers carry everything. 2. **A rolling restart in a fixed order.** Each node in turn sheds what it was serving onto nodes already restarted or not yet touched. The last node in the order is left holding least, and the imbalance survives the restart. 3. **A batch of nodes stopped together**, for a host-level change or a failure domain going away. Everything they served lands on the remainder, and the remainder keeps it. 4. **A node re-added after a long absence.** It rejoins with copies to rebuild and no claim on anything; even after it is fully current it is serving nothing. The common thread is that leadership is *reassigned on loss* and *not reclaimed on return*, so its distribution is a running record of past failures rather than a reflection of the current cluster. ## Why nobody notices Every symptom of drift sits in a per-node view, and every default dashboard is a cluster-wide one. Copies stay even — no records moved — so storage looks correct. Replication looks correct, because every unit still has its copies. What changed is an assignment, and the only number that shows it is the count of units each node currently leads. A cluster can run skewed for months, absorbing it in headroom, and discover it during the first traffic peak or the first loss of the wrong node. | What happened | What moved | What a cluster total shows | |---|---|---| | Node stopped | Leadership for its units | Nothing unusual | | Node returned, caught up | Records it had missed | Brief copy traffic, then nothing | | Nothing hands leadership back | Nothing | Nothing, indefinitely | ## What an interviewer is listening for That you separate the *loss* event, which must reassign leadership immediately or the units stop being served, from the *return* event, which has no reason to reassign anything and therefore does not. Candidates who describe the cluster as self-healing have usually only operated a platform that restores leadership on its own, and have not noticed the machinery doing it for them.

  • A node has been back for an hour and its copies are fully caught up. Why is it still serving almost nothing?
    Because catching up and serving are separate states. Rejoining the caught-up set makes the node eligible to serve, not scheduled to. Until a restoration hands leadership back to the home node — automatically, on a schedule, or because an operator ran it — the units stay where they were reassigned during the outage.
  • Does leadership skew mean the copies are unevenly placed too?
    Not by itself. A cluster can hold a textbook-even placement of copies and still concentrate leadership, because only the assignment moved. The two can also drift together, if nodes were added or evacuated during the same period, which is why the per-node view lists copies held beside units led rather than one number.
  • Is a rolling restart the worst cause of drift?
    It is the most predictable one, not always the worst. A restart in a fixed order reliably leaves the last node serving least, but repeated failures of one node produce sharper concentration because nothing ever gives that node work back. The two compound: a flapping node during a restart sequence is the usual origin story.

saying these in an interview costs you the question

  • A restarted node takes back what it was serving automatically
  • Once copies are caught up, leadership has already returned
  • Skew means the copies were placed unevenly
  • The cluster self-balances leadership, so drift cannot persist
  • Leadership is decided by which node has the most free disk