A cluster node leading a partition (or queue) goes silent: how does the cluster tell a crash from a network cut-off?
answer
- one observation, several causes
- silence does not say which
- a longer timeout only trades errors
- assume it is alive and writing
basics
~20 sThe cluster cannot tell them apart. A crashed node and a node cut off by the network produce the same observation from outside — silence — so the old leader has to be treated as possibly alive and still accepting writes.
solid answer
~50 sFrom the rest of the cluster there is exactly one piece of evidence: the node has stopped answering. A process that died and a healthy process behind a broken link look identical, and waiting longer does not turn silence into proof of death — a longer timeout only changes how often you are wrong and how long the partition (or queue) goes unserved. So an operator never gets to assume "it stopped". The cluster restores service by putting another copy in charge, and at the same time assumes the old leader may still be running, still holding the clients on its side of the break, and still appending writes it believes are valid. Everything else in this area — stamping writes with a generation number, correcting a client that still routes to the old node — exists because that assumption has to be made.
go deeper
Remember the core fact: from outside, a crashed node and a cut-off node look the same — both are silent. Never say "it went down so it stopped writing".
Explain why tuning the unreachability timeout cannot settle it: a shorter one declares healthy nodes unreachable more often, a longer one leaves the partition (or queue) unserved, and neither changes the evidence.
Show the operational consequence: you plan for an old leader that is alive and wrong. Returning nodes are suspects, unacknowledged writes are unknown rather than failed, and runbooks cannot ask you to confirm a node is down.
Frame it as what your standards forbid assuming. The estate's incident language, alert wording and automation must all treat unreachability as a view rather than a state, or someone eventually automates a destructive action on a healthy node.
## What the cluster can actually observe In a broker or streaming platform, a **unit of ownership** — a partition (or queue) — is served at any moment by one **leader** node, while other nodes hold **copies** of its records. When the leader stops answering, the rest of the cluster holds exactly one piece of evidence: *requests to that node are not being answered*. That is a single observation, and several very different situations produce it. | What actually happened | What the rest of the cluster sees | What the silent node sees | What it may still be doing | |---|---|---|---| | The process exited | no answers | nothing; it is gone | nothing | | The host lost power | no answers | nothing | nothing | | A network link between it and its peers broke | no answers | its peers went silent | serving writes to clients that can still reach it | | The node stalled — a long pause, a saturated disk | no answers, then answers again | a jump in elapsed time | resuming, still believing it leads | The rows differ in what is true. The middle column — the only one the cluster can read — is identical in all of them. ## Why waiting longer does not convert silence into knowledge The instinct is to tune the unreachability timeout until it "tells the truth". It cannot, and it is worth being precise about what it does change: - **A short timeout** declares healthy nodes unreachable more often, so ordinary pauses turn into leadership changes and the copy traffic that follows them. - **A long timeout** leaves the partition (or queue) unserved for longer while writers retry. - **Neither setting changes the kind of evidence available.** At the instant the timeout expires, the cluster still does not know which row of the table it is in. - Only the silent node itself could distinguish them, and by construction it cannot say so — the channel that would carry the answer is the thing that may have failed. ## The assumption an operator does not get to make Because the two cases are indistinguishable, "it stopped" is never available as a premise. Practically, the cluster has to hold all of these open at once: 1. Writers on the far side of a break may still be reaching the old leader and getting answers from it. 2. The old leader may still believe it owns the unit and may still be appending records to its local storage. 3. It may come back at any moment carrying records the new leader never saw. 4. Anything it accepted after contact was lost is at best provisional — it has no way to know whether it still has the authority it is exercising. This is the difference between a membership change that is safe and one that quietly loses or duplicates work: the cluster plans for an old leader that is alive and wrong, not for one that is conveniently dead. ## What production designs do instead of deciding Rather than trying to answer an unanswerable question, platforms make the answer unnecessary by making **authority checkable at the receiving end**. Each change of leadership raises a **generation number**, which is stamped on writes and carried in the ownership view clients cache. A node that is no longer the leader can keep trying — its writes carry an older stamp and are refused wherever they have to be accepted by another node. Nobody has to establish whether it is alive first. ## Where platforms differ - On platforms where **one node leads a unit of ownership** and others follow, the cut-off leader is fenced by that stamp when it tries to get a write accepted or acknowledged. - On designs that **commit by majority write**, a node on the minority side of a break simply cannot gather enough acknowledgements, so it stops being able to commit on its own. - On designs with **detached storage**, where records live on shared or remote storage rather than a node's local disk, the refusal happens at the storage layer instead. - **Queue-shaped brokers** are not exempt: whichever node owns the queue can be cut off the same way, even though there is no stream of numbered records behind it. - Some platforms have a node **demote itself** after a period without contact; others rely only on rejection later. A hosted platform chooses for you and may not say which. ## What this changes in the operations room - A runbook step that says "confirm the node is down" is usually not executable; what you can confirm is that you cannot reach it. - Treat a node that comes back as a suspect until it has re-joined under the current generation, not as capacity you immediately reuse. - A write whose acknowledgement never arrived is **unknown**, not failed — the record may exist. - An alert that reads "node unreachable" is a statement about your monitoring's view, not about the node's state, and phrasing it that way saves an incident call from arguing about the wrong thing.
- If the timeout cannot settle the question, what is it actually for?It sets how long the cluster tolerates silence before acting. That is a trade between reacting to brief pauses — each one costing a leadership change and the traffic that follows — and leaving the partition (or queue) unserved while writers retry. It is a tuning choice about cost and responsiveness, never a truth test.
- A node that was cut off comes back with records the new leader never saw. What happens to them?They were accepted under an older generation, so they are not part of the unit's accepted history. The returning node has to discard the divergent tail and re-synchronise from the current leader before it can serve again. If a writer was told those records were acknowledged, that is a durability problem, not a fencing one.
- Does a stalled node behave differently from a genuinely cut-off one?Not from outside — both are silent. It differs afterwards: a stalled node resumes with no sense that time passed and may act immediately on beliefs that are now out of date, which is exactly why authority is re-checked at the receiving end rather than assumed by the sender.
saying these in an interview costs you the question
- Says a node that stopped answering has obviously stopped writing
- Claims a long enough timeout proves the node is dead
- Assumes the cut-off node instantly knows it lost contact
- Treats promotion as proof the old leader released its clients
- Thinks a failed health check reveals which failure occurred