skip to content

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%

answer

  1. who owns the redundancy here?
  2. the store replicates, not the brokers
  3. the number may count serving nodes
  4. mind the span before storage

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.

solid answer

~50 s

Where broker nodes write stream data into a shared durable store - object storage, or a replicated block layer - the number of copies of the *bytes* is a property of that store, chosen by whoever runs it, and not a per-stream dial you set. A broker-side count, if the design exposes one at all, is then about the availability of the serving and leadership path rather than about how many copies exist. Two questions replace "how many copies?": what failure domains does the underlying store replicate across, and is there a span in which an accepted record lives only on one broker node before it reaches the store? That second span is where these designs actually lose data, and it is covered by holding the same un-stored batch on a second node or in a small replicated intermediate tier - not by a number on the stream.

go deeper

for a junior

Know that not every platform lets you choose how many copies of a stream exist; in some, durability comes from storage underneath the brokers and the number you see means something else entirely.

for a middle

Explain who owns redundancy in each design - your per-stream choice, a fixed pair, or a store that replicates internally - and why a broker-side count of one is not automatically alarming.

for a senior

Identify the exposure these designs actually carry: the span between accepting a record and placing it in the store, and how a given design covers it. Then say where the placement decision moved to.

for a principal

Write the standard so it survives the architecture: a requirement in independent failures and domains rather than a copy count, plus a rule for comparing tiers when the platform owns the number.

## Three shapes a copy can take The word "copy" is stable across this product class but the mechanism behind it is not, and a candidate who carries one design's number into another gets the durability posture wrong in a way no configuration review will catch. 1. **N broker-held copies.** The cluster stores the stream's data N times on N nodes, related either as a leader with followers pulling from it or as a set that commits by majority. N is yours to choose, per stream, and placement is a real decision. 2. **A mirrored copy.** Some broker designs keep exactly one paired copy rather than a configurable N. The question is not "how many" but "is mirroring on for this destination", and the placement question shrinks to which node holds the pair. 3. **Shared durable storage underneath.** Broker nodes accept writes and place the bytes into a store that already replicates internally. Nothing per-stream is being copied by the brokers at all. ## What a count means in each | Design | What the number counts | Who owns redundancy | Placement decision | |---|---|---|---| | N broker-held copies | Stored replicas of the data | You, per stream | Which domains the copies span | | A mirrored copy | On or off, one pair | The platform's design | Where the pair sits | | Shared durable store | Serving or leadership nodes | The store underneath | What the store spans, plus node spread | The third row is the trap. A number of one there does not mean the bytes exist once; it usually means one node is responsible for serving a stream whose data is already held several times by the store. Read as though it were the first row, it looks like a catastrophic misconfiguration - and read the other way round, a genuine single-copy stream in the first design looks acceptable. ## The two questions that replace "how many copies?" - **What does the store replicate across, and how much of it do you control?** Usually: several copies inside one zone, or across several zones, with the choice made per storage container rather than per stream. Your per-stream durability decision collapses into which store, container or tier the stream's data lands in. - **What covers a record between acceptance and storage?** A broker that answers a writer before the bytes are in the shared store has created a span in which the record exists only where that node put it. Designs cover this in different ways: holding the same un-stored batch on a second node, writing it first to a small replicated intermediate tier, or simply not answering until the store has it - which costs latency. That second question is the important one and it is easy to miss, because the shared store's own redundancy is advertised loudly and covers only what has arrived in it. ## Where the placement decision goes It does not disappear; it changes level: - you no longer choose which hosts hold copies, but you do choose whether the store's redundancy spans more than one availability zone; - the serving nodes still have a placement of their own - if they are all in one zone, a zone loss is an availability event even when no bytes were lost; - reachability becomes part of durability: data safely held in a store you cannot reach during an incident is not serving anyone; - on a managed offering the choice may be reduced to a tier, in which case comparing tiers *is* the placement decision. ## What a migration gets wrong The recurring failure is a number carried across. A team that standardised on "three copies everywhere" arrives at a shared-storage design, sees a per-stream setting that will not go above one, and either concludes the platform is unsafe or - worse - concludes their standard is met because the word matches. Both are category errors. The standard should have been written as *how many independent failures must a stream survive, and across which domains*, because that question has an answer in all three designs and the number does not. ## What stays true everywhere Durability is still counting independent failures. The unit has moved - from copies on nodes you place, to redundancy inside a store you rent, to a pair you switch on - but the interview answer is the same shape in all three: name what must fail simultaneously to lose the data, and name the span during which fewer things must fail than you think.

  • Where is data most exposed in a design that accepts writes locally and places them into a shared store?
    Between acceptance and storage. Bytes the broker has already answered for but not yet placed live wherever that node put them, so losing the node risks that span. What matters is what covers it - the same batch held on a second node, a small replicated intermediate tier, or waiting for the store and paying the latency.
  • Does such a design still have a placement decision?
    Yes, at a different level. You choose whether the store's redundancy spans more than one availability zone, how the serving nodes are spread, and whether the store stays reachable when a domain is gone. Data that is safe but unreachable during an incident is still an outage.
  • How should a durability standard be written so it survives a move between these designs?
    As a requirement, not a number: how many independent failures a stream must survive and across which domains. That phrasing has an answer whether durability comes from broker-held copies, a mirrored pair or a replicated store, whereas a standard written as a copy count only means something in one of the three.

saying these in an interview costs you the question

  • Assumes a count of one always means one copy of the bytes.
  • Assumes the store's redundancy covers records not yet written there.
  • Carries a familiar copy number to a design with no such dial.
  • Thinks shared storage removes every placement decision.
  • Treats the store's internal copies as a per-stream setting.