skip to content

The leader of a partition (or queue) is lost and no copy in the caught-up set is reachable; what does promoting a copy outside that set, within the same cluster, cost you?

level: seniorimportance: must knowfreq 66%

answer

  1. both options are bad
  2. availability against acknowledged records
  3. the loss reaches nobody as an error
  4. catch-up distance sizes the damage
  5. decide the policy before the incident

basics

~20 s

It costs acknowledged records. A copy outside the caught-up set is missing writes the cluster already confirmed to producers, and promoting it discards them silently — no error reaches anyone. The alternative is keeping every record and staying unwritable for an unknown time.

solid answer

~40 s

This is the availability-against-loss choice, and both options are bad. Waiting for a copy from the caught-up set keeps every acknowledged record, but the partition serves nothing until such a copy returns — which may be minutes, or never if its storage is gone. Promoting a copy outside that set restores writes immediately and throws away every record the cluster acknowledged that this copy never received; the producers were told success and will never be told otherwise. The loss is bounded by that copy's catch-up distance at its last contact with the leader, so a cluster that tracks catch-up distance lets you decide with a number instead of a feeling. This is called lossy promotion, and it should be a pre-declared per-stream policy rather than a judgement made at three in the morning.

go deeper

for a junior

Know that some copies are current and some are behind, and that restoring service with a behind one means the newest records are gone.

for a middle

Explain it in terms of acknowledged records: the producer was told success, the promoted copy never held that record, and no error is ever raised about it afterwards.

for a senior

Show you decide with a number — the candidate copy's catch-up distance — that you record the extent of the loss at promotion time, and that the answer differs per stream.

for a principal

The angle is that the loud cost and the quiet cost are not comparable under pressure, so the decision must be pre-declared, owned by whoever is accountable for the data, and revisited as the stream's value changes.

## The situation One node was serving a **unit of ownership** — a partition (or queue) — and it is gone. Normally the cluster promotes a copy from the **caught-up set**: the subset of copies current enough to take over, which by construction holds every record the cluster acknowledged. The hard case is when no member of that set can be reached — they are on the same failed rack, or they died earlier and nobody noticed, or the only one left is the node that just died. What remains is a copy that is behind. ## The two options, stated honestly | | Wait for a caught-up copy | Promote a copy outside the caught-up set | |---|---|---| | Writes resume | When such a copy returns; possibly never | Immediately | | Acknowledged records | All preserved | Those missing from the promoted copy are discarded | | Who is told | Everyone, loudly — the unit is unavailable | Nobody; producers already received success | | Bound on the damage | Unavailability, unknown duration | The promoted copy's catch-up distance at last contact | | Reversible | Yes, by waiting longer | No, once writes land on the new leader | The second column is often described as 'choosing availability', which is accurate but too comfortable. What is actually being chosen is **silent** loss of records the system already promised to keep. That asymmetry — a loud cost versus a quiet one — is why the fast option wins under pressure and why the decision should not be made under pressure. ## Why the loss is silent Nothing in the system is in a position to complain: - the producers received an acknowledgement at write time and have long since moved on; - the new leader has no knowledge of records it never received, so it cannot report them missing; - readers that had already read past the promoted copy's end find that the stream has become shorter than it was, and the records they saw are no longer there to re-read; - downstream totals computed before and after the event disagree, and nothing explains the difference. The operational consequence is that **you must record the extent yourself**: the end of the promoted copy against the last acknowledged position of the lost leader. If that is not captured at promotion time, nobody can later say what was lost, and every downstream discrepancy becomes an unbounded investigation. ## The number that sizes the decision A behind copy is not a uniform risk. A copy a few hundred records behind and a copy that fell out of the caught-up set an hour ago are different decisions, and **catch-up distance** — how far behind its leader a copy was, in records, bytes or seconds — is what separates them. Practically: 1. Read the candidate's catch-up distance at its last contact with the lost leader. 2. Translate it into the business unit it represents: orders, payments, sensor samples, one minute of telemetry. 3. Compare that against what the same interval of unavailability costs. 4. Decide, and record the number you decided on. A cluster that does not track catch-up distance continuously cannot do step one, and therefore cannot make this decision rationally at the moment it matters. That, not the promotion itself, is the preparation this scenario is really testing. ## Where designs differ - On designs that commit a write only once a **majority** of nodes hold it, any node that is allowed to take over already holds every committed record, so this specific trade does not arise in the same shape. Instead the unit simply stays unwritable while a majority is unreachable. - On platforms with **detached storage**, where the unit's records live on shared or remote storage rather than a node's local disk, promotion copies nothing and a promoted node sees the same records the old one did, so there is no behind copy to promote. - On queue-shaped brokers where a record is removed once it has been handled, the equivalent hazard is not a shortened stream but records that were accepted and are now unrecoverable from any node. - Some platforms let the operator permit this behaviour in advance per stream; others forbid it outright and leave you with only the waiting option. ## Judgement, not a rule The defensible senior answer names the trade in terms of acknowledged records, sizes it with catch-up distance, and then splits by what the stream carries. A stream of metrics that is re-derivable at the source is a different decision from a stream that is the system of record for money, even inside the same cluster. And if the stream can be republished from upstream, the loss becomes a re-publish task rather than a permanent hole — which is why the useful preparatory question is not 'do we allow this' but 'which of our streams could survive it'.

  • After such a promotion, how do you tell downstream teams what they lost?
    Capture the extent at promotion time: the promoted copy's end against the last acknowledged position of the lost leader, and the wall-clock interval that covers. Convert it into the business objects it represents and publish that to every consumer of the stream. If it is not captured at the moment, it usually cannot be reconstructed later, and every discrepancy afterwards becomes unbounded.
  • Would waiting have been the right call if the caught-up copy's storage had actually been destroyed?
    No — then waiting buys nothing and loses everything else. The decision depends on whether a caught-up copy plausibly returns. If its storage is gone, the choice is between an outage of unbounded length and the same records lost anyway, so restoring service and recording the extent is the better call.
  • Does having more copies remove this decision?
    It makes the situation rarer, not impossible: more copies spread across failure domains make it likelier that a caught-up one survives. But copies that share a failure domain fail together, and copies that have fallen behind do not count regardless of how many there are. Placement and catch-up health matter more than raw count.

A bank branch keeps its only current ledger in a safe whose door has jammed. It can stay closed until a locksmith arrives — nobody is served, but every deposit is intact — or it can reopen this morning using last week's duplicate ledger. Reopening restores service instantly, and every deposit made since that duplicate was written simply does not exist; the customers hold stamped receipts that nothing will honour, and nobody is notified that this happened.

saying these in an interview costs you the question

  • Calls promoting a behind copy safe because service came back
  • Thinks producers are notified when their acknowledged records vanish
  • Believes the new leader can recover what it never received
  • Says more copies would have made the promotion lossless regardless
  • Treats the choice as cluster-wide rather than per stream
  • Cannot bound the loss because catch-up distance was never tracked