skip to content

For a week a stream's caught-up set has held only its leader, yet no write ever failed — what has actually been true of those writes?

level: seniorimportance: should knowfreq 50%

answer

  1. the promise is relative to a set
  2. a set of one satisfies it trivially
  3. no error, no gap, one machine
  4. rejoining needs surplus rate, not repair

basics

~20 s

Every write in that week was held by exactly one machine. A promise defined against the caught-up set is satisfied trivially once that set is the leader alone, so the writes look acknowledged and safe while nothing else holds them.

solid answer

~50 s

The durability promise is expressed **relative to the caught-up set**, so when the set narrows to the leader alone the promise is still met — by one machine. For a week, records the writers believed were replicated existed in exactly one place, and losing that leader's volume would have lost them. Nothing failed loudly because the condition is self-satisfying: there was no rejected write, no error and no gap in the stream, just a set that quietly stopped containing anyone but the leader. The usual cause is that the followers dropped out for an ordinary reason — a slow disk, a saturated link, a restart that left them far behind — and then never had enough spare throughput to close the gap while the leader kept writing. The lesson is that the width of the set, not the configured copy count, is what a durability claim is measured against.

go deeper

for a junior

Take away one idea: a rule that says "every current copy has it" is satisfied even when only one copy is current. Success answers do not by themselves prove a record is on several machines.

for a middle

Explain how a set narrows without anyone acting, and why a relative promise still returns success at width one. Be able to say what a hardware loss would have cost that week.

for a senior

Demonstrate the diagnosis and the arithmetic of recovery: a copy rejoins only with throughput above the leader's incoming rate, so the fix is relieving the bottleneck, not restarting the process.

for a principal

Own the estate question this raises — what standard makes an unintended posture visible on streams nobody reviews, and what each durability tier commits to when the set narrows.

## The condition A stream is configured with three copies. At some point in the week, both followers fell out of **the caught-up set** and never returned. The leader kept accepting writes, writers kept receiving successful answers, readers kept reading, and every graph of throughput looked exactly as it did the week before. What changed was invisible from the writer's side: for seven days, the stream's data existed on **one machine**. ## Why nothing failed This is the trap worth understanding, because it is a property of how the promise is phrased, not a bug. A durability requirement of the form *the write counts once every currently caught-up copy holds it* is **relative**. It quantifies over a set whose membership moves. When the set is three copies wide, the requirement is strong. When the set is the leader alone, the requirement reads *the write counts once the leader holds it* — which is satisfied the instant the leader has written it. The rule was honoured on every single write; it had simply stopped meaning anything. So the failure mode is: - no rejected writes; - no errors surfaced to any writer; - no gap or reordering in the stream; - no reader complaint, because reading is unaffected; - and a week of records with a single point of failure under them. Some platforms offer a separate rule that refuses writes unless a minimum number of copies is currently caught up; where such a rule exists and is set, this silent condition becomes a loud one instead — an availability event rather than a durability one. Where it is absent or left at its most permissive value, the condition stays silent. ## What was actually true of the week's writes | The writers believed | What was true | |---|---| | The record is on several machines | The record is on the leader's volume only | | A node loss is survivable | A node or volume loss takes every record written that week that the followers never pulled | | The copy count is the durability posture | The copy count was unchanged and irrelevant; the set's width was one | | A successful answer means replicated | A successful answer means the stated rule was met, whatever it currently quantified over | Note what is *not* on the right-hand column: the records were not corrupted, the stream was not inconsistent, and readers were not served wrong data. The exposure is entirely about what a single hardware failure would have cost. ## Why the copies never came back Dropping out is ordinary; staying out for a week is the part that needs explaining. The usual reasons are all about **rate**, not about breakage: 1. **The copy has no surplus throughput.** To rejoin, a follower must read fast enough to close its accumulated gap *and* keep pace with everything the leader is still writing. A copy whose disk or link runs at roughly the incoming rate never closes the gap; it merely holds it steady, forever outside the boundary. 2. **The gap grew during a period of high load** — a backfill, a replay, a seasonal peak — and the load never fully receded, so the surplus needed to catch up never appeared. 3. **The copy shares hardware with something noisy**, so its effective read rate is far below the leader's write rate whenever that neighbour is busy. 4. **The copy was restarted far behind** and the cluster's catch-up traffic is deliberately paced, so it recovers slowly by design. In each case the copy keeps pulling and keeps holding old records; it is doing work, just not enough of it. ## Getting back to a wide set The route back is arithmetic, not a toggle: the lagging copy needs more read-and-write throughput than the leader's incoming rate, at least for as long as the gap takes to close. In practice that means relieving whatever caps it — moving the copy off contended hardware or a saturated path, giving catch-up traffic more headroom during a quiet period, or reducing what the leader is ingesting for a while. A copy relocated to fresh hardware starts from an empty gap but must transfer the stream's retained data first, which is its own cost. ## Where the shape is different This specific silent narrowing is a property of designs that keep a standing membership of current copies and express the durability promise against it. On platforms that commit a write only when a majority of copies has answered, the write itself stops succeeding when too few copies can answer — the condition surfaces as failed or slow writes rather than as a quiet one-machine posture. And in designs where durability comes from shared underlying storage, there is no per-node caught-up set to narrow. The transferable lesson is the general one: know what your platform's success answer actually quantifies over.

  • Why is this described as a durability problem rather than an availability problem?
    Because nothing stopped. Writers were served, readers were served, and the stream stayed intact — availability was perfect all week. What was degraded was how much hardware loss the data could survive, which is a durability property and shows up only when the loss actually happens.
  • A copy has been outside the set for days and the operator restarts its node. Does that help?
    Usually not on its own. A restart clears a stuck process but does not change the arithmetic that keeps the copy behind: it still has to read faster than the leader is writing to close the gap. If the cause is a saturated link, contended hardware or a sustained ingest rate above the copy's throughput, the restart just resets the same race, often with a larger gap.
  • The copy count still reads three the whole time. Why is that number so misleading here?
    Because it answers a question about storage placement, not about currency. Three copies existed the entire week and two of them held only week-old records. The count is what was asked for; the width of the caught-up set is what was delivered, and only the second one bounds what a hardware loss would cost.

saying these in an interview costs you the question

  • Says the writes were fine because no error was returned
  • Believes a promise stated over the caught-up set cannot be met by one machine
  • Assumes a lagging copy will catch up on its own once its hardware is healthy
  • Reads the configured copy count as the durability the stream actually had
  • Calls this an availability incident because nothing stopped serving
  • Thinks the records on the lagging copies were corrupted rather than merely old