How long should the cluster tolerate silence from one of its nodes before promoting new leaders, and what does shortening that tolerance cost?
answer
- silence is not the same as death
- both directions of the dial hurt
- ordinary pauses set the floor
- needless promotion still opens a gap
- pair it with the writer's budget
basics
~20 sTolerance must exceed the node's ordinary pauses — garbage collection, a slow disk flush, a congested link — or the cluster promotes replacements for nodes that are alive. Shortening it shortens outages but buys needless promotions, each with its own gap and client churn.
solid answer
~50 sThe dial is how long a cluster node may stay silent before the cluster treats it as gone and promotes replacements for the partitions (or queues) it led. Set it too long and every real node death costs that much unavailability before anything starts. Set it below the node's ordinary stalls — a long pause inside the process, a disk flush that blocks, a brief network problem — and the cluster promotes replacements for a node that never died, and each needless promotion has its own gap in which no node serves writes, plus clients refreshing stale owner maps and leadership landing where it was not planned. The honest method is to set the tolerance from the observed distribution of real pauses with headroom, not from a round number, and to size the writers' retry budget against the resulting worst case.
go deeper
Know that a node going quiet is not proof it died, and that the cluster waits a set amount of time before replacing the leaders it served.
Explain both directions: too long means every real failure costs that wait, too short means replacements are promoted for nodes that merely paused, each with its own outage and client churn.
Show that you set the value from the measured tail of real pauses in that specific environment, and that you check the resulting worst case against what the writers can survive.
The angle is that this is a cost allocation, not a tuning exercise: a short tolerance spends reliability to buy recovery speed, and on a hosted platform you may only own the client half of the pair.
## What the dial actually controls A cluster has to decide when a silent node is a dead node. The tolerance is the interval of silence it allows before it declares the node gone and promotes replacements for every **unit of ownership** — every partition (or queue) — that node was leading. It is one of the few knobs whose two failure directions are both bad and both common, which is why interviewers use it. The important framing is that the tolerance is not a detection quality setting. Detection is not made more accurate by shortening it; it is made **faster and less certain**. What the operator is really choosing is the ratio between two costs. ## The cost of tolerating too much silence Every real node death now costs at least that interval before anything begins. During it: - the units that node led have no leader, so writes to them are refused; - the writers' retry budgets are burning against a gap nobody has started to close; - the cluster looks healthy by node count, because the node has not yet been declared gone. A long tolerance is therefore paid for on exactly the occasions it was meant to protect — genuine failures — while it buys nothing on the days nothing fails. ## The cost of tolerating too little A cluster node goes quiet for many reasons that are not death: a long pause inside its process, a disk flush that blocks for seconds, a link that saturates, a host that is briefly starved of processor time. If the tolerance is below the tail of those events, the cluster promotes replacements for a node that is alive and about to answer. Each such promotion is not free: - it opens its own gap in which no node serves writes for those units; - every client holding a stale owner map has to discover the change and refresh; - leadership ends up on nodes the plan did not intend, and has to be corrected later; - if the copies that were current are themselves not ideal candidates, a needless death forces a real promotion decision that nobody wanted to make today. The pattern to recognise in an incident review is a cluster that promotes replacements for the same node repeatedly while that node's own records show it never stopped running. The fix is the tolerance, not the node. ## Setting it honestly | Input | What it tells you | |---|---| | Observed distribution of node pauses, at the tail rather than the median | The floor the tolerance must clear | | Worst normal network hiccup between nodes, including across failure domains | Why a stretched cluster needs more tolerance than a single-room one | | How much write unavailability each stream can absorb | The ceiling you are willing to pay on a real death | | The writers' total retry budget | Whether the worst case is survivable by the applications at all | The method is: measure the pauses you actually have, put the tolerance above their tail with headroom, then check that the resulting worst-case gap still fits inside the writers' budget. If it does not, one of the two has to move — and raising the client budget is usually cheaper than shortening the tolerance, because it costs latency rather than correctness. ## The matching dial on the client The second half of the pair belongs to the writer: how long it keeps retrying before reporting failure upward. These two are meaningless apart. A tolerance of tens of seconds with writers that give up in a few is a design that turns every node failure into application errors; writers that retry for minutes against a cluster that promotes in seconds simply never notice failures at all. Tune them as one number with a margin between them. ## Where designs differ How the silence is noticed varies more than the choice it drives. Some clusters keep their own membership records internally; others depend on a separate coordination role to notice and record the change; hosted offerings often do not expose the tolerance at all, so the interval is the provider's and the only dial you hold is the client's budget. On designs with **detached storage**, where records live on shared or remote storage, promotion itself is cheap, which makes a shorter tolerance more affordable — but it does not make a promotion triggered by a pause harmless, because the churn on clients is unchanged. ## What to say in an interview Name both directions, say that the tolerance must clear the tail of ordinary pauses rather than the average, and finish on the pairing with the writer's budget. Candidates who give only one direction — usually 'make it short so failover is fast' — are the ones who have never watched a cluster shuffle leadership all night because a node paused for four seconds at a time.
- Your cluster promotes replacements for the same node twice a week, and that node's own records show it never stopped. What do you change?Raise the tolerated silence above the tail of that node's real pauses, after measuring them rather than guessing. Look for the underlying stall too — a blocking flush, a saturated link, a starved host — because the tolerance only hides it. Restarting or replacing the node fixes nothing if the tolerance is still below its normal behaviour.
- Why can a tolerance that works in one cluster be wrong in another running the same software?Because it is set against the environment, not the product. A cluster spread across failure domains has longer and more variable round trips than one in a single room; a host with noisy neighbours pauses differently from a dedicated one; and larger stored state lengthens stalls. The correct value is the tail of that cluster's own pauses plus headroom.
- If you cannot change the tolerance at all, what is left to tune?The writer's side. Size the total budget before failure is reported upward so it outlasts the worst gap the platform can produce, and make sure the client's buffer can hold the records accumulated in that interval. That converts an unavoidable gap into latency rather than errors, which is usually the cheaper cost.
saying these in an interview costs you the question
- Says shorter detection is strictly better because failover is faster
- Assumes a silent node has stopped doing work
- Treats a promotion caused by a pause as harmless
- Picks a round number instead of measuring real pauses
- Tunes the cluster tolerance without looking at the writers' retry budget
- Blames the node for repeated promotions when the tolerance is the cause