skip to content

An unreachable host's replicas were replaced; the host returns with its old copies still running and writing - what did the platform actually guarantee?

level: seniorimportance: should knowfreq 42%

answer

  1. silence is not proof
  2. a floor, not a ceiling
  3. availability chosen over exclusivity
  4. the far side stops its own copies
  5. bound writers at the data, not the scheduler

basics

~20 s

A floor, not a ceiling. The platform promises that at least the declared number of copies exists; it cannot promise at most that number, because declaring a silent host's workloads gone is a guess made without contact. Bounding writers must come from outside the scheduler.

solid answer

~50 s

When a host stops reporting, the control plane cannot distinguish a dead machine from a partitioned one, and it chooses availability: after the timeout it writes the copies off and starts replacements. If the truth was a partition, the originals kept running the whole time, and for that window the workload had more copies than declared, all doing real work. What ends the overlap is reconciliation on the host's side - once its agent can talk to the control plane again, it stops what it is no longer assigned to run. Nothing can end it during the partition. So if two simultaneous writers would corrupt something, the bound has to be enforced where the data is: a lease that expires, a token the store rejects when it is stale, or storage that only one instance can attach at a time.

go deeper

for a junior

Take away the core fact: the platform makes sure enough copies exist, not that too many never do. A host it cannot reach may still be running the copies it wrote off.

for a middle

Explain why the control plane cannot tell a dead host from a partitioned one, and why it chooses to replace anyway rather than wait for proof that may never arrive.

for a senior

Describe the overlap window end to end, including that the stop is executed by the returning host's own agent, and name the data-layer mechanisms that bound writers when duplication is unsafe.

for a principal

Frame it as an availability-versus-exclusivity choice made once, cluster-wide, and decide which workloads must buy exclusivity at the data layer and what being short of copies during a partition costs the business.

## The guess at the centre of it Self-healing on an unreachable host rests on an inference the platform cannot verify: *no reports means the workloads are gone*. The deciding half has no contact with the host, which is the entire premise, so it also has no way to confirm the conclusion or to act on it - it cannot stop what it cannot reach. When the timeout expires it does the only thing available: it stops counting those copies and lets the replica loop start new ones elsewhere. If the machine really was dead, that is exactly right. If the network was broken and the machine was fine, the ingester copies on it carried on reading from their sources and writing their output for the entire partition, and now there are more copies than the spec declares. ## Floor, not ceiling This is the precise statement of what a replica count buys you: | The platform does promise | The platform does not promise | |---|---| | To notice when fewer copies exist than declared | That no more than the declared number ever run | | To keep starting copies until the declared number is reached | That a copy it stopped counting has actually stopped | | To stop surplus copies once it can reach them again | Anything at all about a host it cannot reach | The choice is deliberate, and it is a choice for **availability**. A platform that refused to replace until it could prove the originals had stopped would leave workloads short for the entire length of any partition - which is unbounded, since a partition that never heals never proves anything. For interchangeable replicas of a stateless service, running six copies for a few minutes instead of five is harmless, and running three instead of five is an outage. The design picks the harmless failure. ## How the overlap actually ends 1. The host regains contact and its agent re-registers with the control plane. 2. The agent compares what it is currently running against what the control plane says is assigned to it. 3. The copies that are no longer assigned are stopped by the agent, on the host, locally. 4. The workload settles back to the declared count. Note which side performs step 3. The stop is executed by the half that can reach the containers, driven by the state the deciding half holds. That is why nothing can shorten the overlap from the control plane's side: the only actor able to end it is on the far side of the break. A host that never comes back never has its copies stopped by anyone - they end when the machine does. ## Where a real bound has to live If two concurrent instances of this ingester would double-write, corrupt a file, or both act as the single owner of something, the orchestrator is the wrong layer to ask. The bound belongs where the shared resource is: - **A lease with an expiry.** An instance holds a time-bounded claim and must keep renewing it. A partitioned instance cannot renew, so it loses the claim and must stop acting on it before it expires - which makes the safety a property of its own clock, not of anyone reaching it. - **A token the store validates.** Each new holder gets a strictly increasing token; the store rejects writes carrying an older one. A stale instance's writes fail at the point of impact even though it believes it is still in charge. - **Exclusive attachment.** Storage that admits one writer at a time refuses the second attach, so the replacement either waits or fails to start rather than joining the original. - **Work that tolerates duplication.** Making the write idempotent, or keyed so a repeat overwrites itself harmlessly, removes the need for a bound at all - usually the cheapest answer for an ingest path. Each of those is a property of the data path, deliberately outside the scheduler. Note the trade they all share: guaranteeing at most one writer means accepting periods with none. ## What an interview is checking Candidates who have only read about self-healing describe it as a guarantee that the declared number of copies is running. Candidates who have operated through a partition describe it as a guarantee about the lower bound, with an upper bound that is best-effort and briefly violated by design. The second answer is the one that predicts the incident: duplicated readings, a file written by two processes, or a scheduled task that fired twice, all during a window when the dashboard said the workload was healthy at its declared count the whole time.

  • Could the platform simply stop the old copies when it declares them gone?
    Not while the host is unreachable - that is the premise of the situation. Any stop command has to travel the same path that the reports stopped travelling on. The write-off is a bookkeeping change on the deciding side, and the actual stop happens later, executed locally by the host's agent once contact returns.
  • How does an ingest path make this overlap a non-event?
    By making duplicate work harmless: key each reading so a repeat lands on the same record rather than appending a second one, and deduplicate on the reading's own identity rather than on arrival order. Then two copies producing the same output for a few minutes costs some wasted work and changes nothing downstream.
  • Does a workload with stable per-instance identity change this picture?
    It changes the platform's default posture: where an instance is named rather than interchangeable, platforms are typically much more reluctant to start a replacement for an unreachable one, precisely because two instances claiming one identity is unsafe. The trade is the same one stated the other way - such a workload stays short until someone confirms the original is gone.

saying these in an interview costs you the question

  • Assumes the platform guarantees only one copy of a replica ever runs
  • Says an unreachable host is by definition a dead host
  • Thinks fast replacement is free because the old copies must be gone
  • Assumes the scheduler alone prevents two instances writing the same data
  • Believes the overlap ends when the control plane declares the copies gone