skip to content

Copy Sets & Durability

How many copies of a record exist, which must accept a write before it counts, and what happens when that set shrinks. Asked because lost records trace to a durability setting nobody revisited.

part ofBroker & streaming operationsoverview, primer and where to startread it →
on this pageshow

questions

18

Before a broker answers a write, what can the writer be made to wait for, and what does each option cost in write latency?

level: juniorimportance: must knowfreq 72%

answer

  1. a dial, not a guarantee
  2. four rungs of waiting
  3. each rung, one more holder
  4. top rung pays the slowest copy
  5. chosen per stream

basics

~20 s

A write can be answered with no wait at all, once the leader holds it, once a majority of copies hold it, or once every caught-up copy holds it. Each rung up adds a network round trip of latency and removes one way to lose the record.

solid answer

~50 s

The acknowledgement rule is the server-side requirement met before a writer is told its record was accepted. The common ladder has four rungs: no wait (fire-and-forget), the leader alone, a majority of copies, and every caught-up copy. Fire-and-forget is the fastest and loses records to any failure or even a dropped connection, because nothing confirmed anything. The leader alone costs one hop and survives nothing worse than losing that leader. A majority or every caught-up copy costs the round trip to the slowest copy that has to answer, and in exchange the record survives losing the leader. The rule is a durability-versus-latency dial, and it is chosen per stream, not once for a whole cluster. Note that accepting a record is not the same as forcing its bytes to persistent media — that is a separate step.

go deeper

for a junior

Be able to name the rungs in order and say, for each, how many places hold the record when the writer is told it was accepted. That alone answers the screening version of this question.

for a middle

Explain where the latency of each rung comes from: a hop for the lower rungs, the slowest participating copy for the top one. Add that the rule is set per stream on top of a cluster default.

for a senior

Show you have chosen this dial under pressure — which streams you put on which rung, and how you noticed a sick machine through producer latency rather than through an availability alarm.

for a principal

Frame it as a cost envelope for the estate: each rung has a latency profile and an exposure profile, and the organisation needs a small number of named combinations rather than one argument per stream.

## What the answer to a write actually means A writer sends a record and the broker eventually answers. **The acknowledgement rule** is the server-side requirement that must be met before that answer goes back. It is the knob that decides how much of the cluster has seen a record at the moment the writer is told "accepted", and therefore how much has to fail before the record is gone. Two things it is *not*. It is not a statement that the bytes have been forced onto persistent media — acceptance into a broker process and a forced write to durable storage are separate steps with separate settings. And it is not about what a consumer later does with the record; a reader acknowledging work it has finished is a different mechanism entirely, on the other side of the system. ## The rungs of the ladder Four rungs is the common shape; some designs expose fewer. - **No wait (fire-and-forget).** The writer sends and carries on. There is no answer to wait for, so nothing can be retried intelligently: a record lost in the network, at a full buffer, or on a broker that was restarting simply never existed. Used for high-volume telemetry where a missing sample costs nothing. - **The leader alone.** On designs where one copy leads a stream and the others follow, the leader writes the record and answers immediately. One network hop of latency. The record exists in exactly one place; losing that one machine before a follower has fetched it loses the record. - **A majority of copies.** The write is answered once more than half the copies hold it. Latency is set by the *median* copy, so one slow machine does not stall writes. This is the natural rule on designs that replicate by majority write rather than by a leader-plus-follower arrangement. - **Every caught-up copy.** The write is answered once every copy the leader currently counts as current holds it. Latency is set by the **slowest** copy that has to answer, which is why this rung is the one that exposes a single sick disk as a producer-side latency spike. ## What each rung buys and costs | What the writer waits for | Holds the record when the answer arrives | Latency driver | Lost by | |---|---|---|---| | No wait | Unknown — possibly nothing | None | Any failure, silently | | The leader alone | One copy | One hop to the leader | Losing that leader before a follower copies it | | A majority of copies | More than half | The median copy | Losing more than half at once | | Every caught-up copy | Every current copy | The slowest current copy | Losing every current copy at once | The important shape of that table: latency rises roughly one round trip from the first rung to the second, and then as a function of *tail* behaviour, not average behaviour, from the third rung up. A cluster whose copies are spread across failure domains pays the inter-domain round trip on the upper rungs. ## The choice is per stream A cluster carries a default, and individual streams take a **per-stream override** on top of it. That matters because durability requirements are a property of the data, not of the machines: a payment event and a page-view event can live on the same cluster with different rules. An estate that sets one rule cluster-wide is either paying tail latency on data that does not need it, or running valuable data at the exposure level chosen for telemetry. ## Designs that do not offer all four rungs This is where a candidate who has only operated one platform gets caught. 1. Where a broker keeps **a single mirrored copy** rather than a configurable number, the choice collapses to two states: answered by the primary, or answered once the mirrored copy also holds it. There is no majority to wait for. 2. Where durability comes from **shared durable storage underneath the brokers** rather than from broker-held copies, the wait is for that underlying store to confirm, and copy count is not the operator's dial at all. 3. Where records are **deleted on acknowledgement** rather than retained in a log, the same ladder still applies to the write, even though there is no stored reading position anywhere in the system. ## What an interviewer is listening for That you state the rule as a trade, name what physically holds the record at each rung, and know where latency actually comes from — a hop for the low rungs, the slowest participating copy for the top one. Reciting rung names without the failure each one still permits is the shallow answer.

  • Why does waiting for every caught-up copy make write latency sensitive to a single slow disk?
    Because the answer cannot go back until the last required copy has the record, the rule takes the maximum over the participating copies rather than the median. One machine with a degraded disk therefore sets the latency for every write on that stream, which is why this rung surfaces hardware trouble as a producer-side latency spike long before anything is unavailable.
  • Does a write that has met the acknowledgement rule mean the record survives a power cut?
    Not by itself. Meeting the rule means the required copies accepted the record into their own storage path; whether those bytes have been forced out of volatile memory onto persistent media is a separate setting. The two questions are deliberately distinct, and an operator who treats them as one is describing a stronger guarantee than the cluster gives.
  • Where is this rule chosen — by the writer or by the server?
    Usually both have a say. The writer asks for a level of waiting, and the server carries defaults plus its own floors that can make a request more demanding but not less. That split matters: an operator who tightens only the server side may find writers still asking for the weakest rung and getting it.

Posting a parcel: drop it in the box and walk away, get a receipt at the counter, get a receipt once it reaches the depot, or wait until every branch on the route has scanned it. Each extra scan takes longer and removes one way to lose it.

saying these in an interview costs you the question

  • Thinks an acknowledged write is already on persistent media
  • Believes fire-and-forget still retries a failed send
  • Assumes the strongest rule is free because copies answer in parallel
  • Treats the rule as one cluster-wide setting with no per-stream override
  • Confuses this server-side wait with a consumer acknowledging finished work
  • Cannot say what physically holds the record at each rung
open as a page

Within one cluster, a stream's copy count goes from two to three - what does the third copy cost in stored bytes and traffic?

level: juniorimportance: must knowfreq 60%

basics

~20 s

The third copy adds a complete extra set of the stream's stored bytes - about fifty per cent more disk, since the base was two - plus one more transfer of every incoming byte, continuously, and a one-off backfill of everything already retained.

open as a page

A broker acknowledges a record as stored, yet a power cut on that machine loses it — why?

level: juniorimportance: must knowfreq 62%

basics

~20 s

An acknowledgement usually means the bytes were accepted into the operating system's file cache, which is volatile memory, not that they were forced onto persistent media. Everything accepted but not yet forced is the loss window a power cut takes.

open as a page

A stream keeps three copies within a single cluster. Which failures does that survive, and which does it not?

level: juniorimportance: must knowfreq 72%

basics

~20 s

Copies survive the loss of whatever they do not share: a broker process, a host, a volume. They do not survive what replication reproduces faithfully - a deleted stream, a bad record, expiry, or losing the whole cluster.

open as a page

A stream's acknowledgement rule is every caught-up copy and its minimum-copies floor equals its copy count; what happens when one copy drops out?

level: middleimportance: must knowfreq 62%

basics

~20 s

Writes to that stream are refused until a third copy is caught up again. Setting the floor equal to the copy count turns the loss of any one copy into a write outage, which is why the floor is normally set one below the count.

open as a page

A stream ingests 400 GB a day, is kept for seven days, and the cluster keeps three copies - what stored-byte figure should planning use?

level: middleimportance: must knowfreq 56%

basics

~20 s

Multiply rate by window by copies: 400 GB a day times seven days times three copies is about 8.4 TB of stream data, then add headroom for per-record overhead, partial segments, peak days and rebuilding a copy - so provision well above the raw product.

open as a page

All three copies of a stream sit on the same host in one cluster. What does the copy count still protect against?

level: middleimportance: must knowfreq 64%

basics

~20 s

Very little: a single broker process failing, and one volume failing if the copies are on separate volumes. Host, storage-array, rack and power failures take all three together, so the count bought almost no independence.

open as a page

How does a cluster decide that a follower copy has fallen far enough behind its leader to drop out of the caught-up set?

level: middleimportance: must knowfreq 64%

basics

~20 s

An allowance called the catch-up window sets how far behind a follower copy may be — a number of records, or elapsed time since it was last current. Past it the leader drops the copy; back inside it, the copy is re-admitted.

open as a page

Within one cluster a stream has three copies, but the leader counts only two as caught up — what does that distinction mean?

level: juniorimportance: should knowfreq 55%

basics

~20 s

Copy count says how many copies exist; the caught-up set says how many are current enough right now to stand behind the leader. A copy that has fallen behind still exists and still holds its older records, but does not count.

open as a page

Why do most brokers force writes to persistent media in batches or on a timer rather than per record?

level: middleimportance: should knowfreq 50%

basics

~20 s

Forcing costs a round trip to the device that cannot be shared between records, so doing it per record collapses throughput to the device's operation rate. Batching amortises one force over many records and leaves a bounded loss window instead.

open as a page

A cluster's default copy count was raised to three, yet a stream created earlier still keeps one copy. Why?

level: middleimportance: should knowfreq 55%

basics

~20 s

Because the copy count is a property of each stream, fixed when that stream was created from whatever default applied at the time. A cluster default is a template consulted at creation, not a rule re-applied to streams that already exist.

open as a page

An operator floors caught-up copies at two on every stream, yet a leader failure still loses acknowledged writes; why did the floor never take effect?

level: seniorimportance: should knowfreq 46%

basics

~20 s

Because the floor is only consulted on a write that waits for every caught-up copy. Writers asking to be answered by the leader alone are answered by the leader alone, and the server-side floor never enters the decision.

open as a page

A durability standard fixes the copy count and acknowledgement rule on a cluster whose storage bill has doubled - which levers cut the total, and in what order?

level: seniorimportance: should knowfreq 44%

basics

~20 s

Measure per-stream first, then act on the multiplicand before the multipliers: fewer bytes per record, then retention windows nobody chose, then closed history moved to cheaper object storage, then dead streams retired. Changing the copy count is a posture decision, not a cost lever.

open as a page

Records were acknowledged by three caught-up copies in one cluster, then a shared power failure lost them — why?

level: seniorimportance: should knowfreq 46%

basics

~20 s

All three copies held the record in volatile memory rather than on persistent media, and one event took the three machines together, so it took the unforced span on each of them. Copies cover machines failing separately; forcing covers them failing together.

open as a page

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%

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.

open as a page

In a design where durability comes from a shared replicated store rather than broker-held copies, what does a copy count mean?

level: seniorimportance: nice to knowfreq 34%

basics

~20 s

Usually not what it means elsewhere. Redundancy belongs to the store underneath, across its own failure domains, and any broker-side number counts serving nodes rather than copies of the bytes. Ask what the store spans, and what covers records not yet in it.

open as a page

A candidate says the caught-up set is always a leader's list of current followers — on which cluster designs is that wrong?

level: seniorimportance: nice to knowfreq 36%

basics

~20 s

Only leader-follower designs keep a standing membership list. Majority-write designs decide eligibility per write — whichever copies answer first — with no list and no drop-out event, and shared-storage designs have no per-node copies to be current at all.

open as a page

How would you define a small set of estate-wide durability tiers so teams pick a write's acknowledgement rule without re-deciding it per stream?

level: principalimportance: nice to knowfreq 34%

basics

~20 s

Publish two or three named tiers, each fixing a copy count, a floor on caught-up copies, an acknowledgement rule and a latency budget, then enforce the choice at stream creation so an unreviewed stream lands on a safe default rather than on nothing.

open as a page