Why is a cross-cluster copy of a stream always behind the source cluster that feeds it?
answer
- two writes, not one
- the copier is a client of both
- acknowledged before it travels
- the span has a name
basics
~20 sA cross-cluster copier is an ordinary client of both clusters: the source stores a record and answers the writer before the copier has even read it. The record therefore exists on the source first and on the target some time later.
solid answer
~50 sThe hop is asynchronous by construction. A writer sends a record to the source cluster, the source stores it under its own durability rule and acknowledges the writer; only afterwards does a cross-cluster copier read that record and append it to a stream on the target cluster. The copier sits outside the source's write path, so nothing on the source waits for the far side. The span between 'acknowledged on the source' and 'stored on the target' is `copy lag between the two clusters`, and it is a property you measure and budget, not a fault you close. Faster links, bigger batches and more copier workers shrink it; nothing sets it to zero, because the writer was already told the write succeeded. The only arrangement without that span is one where the write itself waits for the far cluster, which is a different design with a different bill.
go deeper
Remember the order of events: the source stores the record and answers the writer, and only afterwards does a separate process carry it to the other cluster. That order is why the second cluster is always slightly behind.
Be able to name what the gap is made of - read cadence, distance, the target's own write, copier capacity, deliberate throttling - and to say which of those you can change and which you cannot.
Show that you treat copy lag between the two clusters as a standing operational number with a normal range, and that you know what a growing gap means: the hop is draining slower than writers are producing.
The judgment call is what the organisation pays for: inter-site bandwidth for a narrower gap, against accepting a wider one. Framing a copy hop as something that can be made lossless is how that conversation goes wrong.
## What a cross-cluster copier is A **cross-cluster copier** is a process that reads records from a stream on one cluster - the **source cluster** - and appends those same records to a stream on another cluster - the **target cluster**. One directed source-to-target link like that is a **copy hop**. Some platforms build the copier into the product; others leave you to run it beside the clusters. Either way it behaves as an ordinary client: it reads through the same interface any reader uses and writes through the same interface any writer uses. It is **not part of the source cluster's write path**, and the source does not wait on it. That single fact is the whole answer. Everything below is its consequences. ## Follow one record 1. A writer sends a record to the source cluster. The source stores it according to the durability rule configured **there** and answers the writer. 2. The record now exists on the source cluster and nowhere else - and the writer has already been told it succeeded. 3. Some time later the copier reads that record from the source. 4. The copier appends it to the stream on the target cluster and waits for the target's own acknowledgement. 5. Only now does the record exist on both clusters. | moment | on the source | on the target | what the writer has been told | |---|---|---|---| | after step 1 | yes | no | stored | | during steps 3-4 | yes | not yet | stored | | after step 5 | yes | yes | stored | The span between step 2 and step 5 is **copy lag between the two clusters**. It is normally reported in seconds, in records, or both. ## What that span is made of - **Read cadence.** The copier fetches in batches; a record waits until the batch that includes it is fetched. - **Distance.** A long-distance link adds a round trip in each direction, and that floor is physics, not configuration. - **The target's own write.** The target cluster applies its own durability rule before acknowledging the copier, so the hop inherits that wait too. - **Copier capacity.** How many workers the copier runs against how many parts of the stream, against the rate writers are producing. If the source produces faster than the hop drains, the gap grows without bound. - **Deliberate throttling.** Hops are often rate-limited so the copy does not consume the whole inter-site link; that trades a wider gap for a predictable bandwidth bill. - **Interruptions.** Every copier restart, redeploy or network blip leaves a backlog it then has to work through at above the production rate. ## Why it cannot be driven to zero Each item above can be reduced. None of them can be removed, because the ordering of events is fixed: the source answered the writer *before* the copier held the record. An engineer who promises a zero-lag copy is describing a different arrangement - one in which the write is not acknowledged until the far cluster also holds it. That puts the inter-site round trip inside every single write and makes the source's ability to accept writes depend on the far side being reachable. That is a legitimate design, but it is not a copy hop, and it is not what 'mirroring to a second cluster' means. ## What follows operationally - At any instant there is an **uncopied tail**: records that exist on the source and have not reached the target. If the source is lost right now, those records are not on the target. - Copy lag between the two clusters is therefore a **standing measurement**, watched the way you watch a queue depth - not an alarm that something is broken every time it is non-zero. - Because the gap is real, any claim that the two clusters hold the same stream is only ever true 'as of a moment slightly in the past'. - Turning that gap into a number the business will sign for is a separate exercise from measuring it. ## Where platforms differ The asymmetry is universal in this class, but the details are not. On platforms that split a stream into parts and give each record a numeric position, the copier tracks how far it has read **per part**, and copy lag is naturally a per-part number. On queue-shaped platforms where a record is removed once a consumer acknowledges it, the copier is simply a consumer that republishes elsewhere, and the same asymmetry shows up as 'the message was taken from here and put there afterwards'. Some products ship a copier and some expect you to supply one; some let the copy run at line rate and some cap it by default. What never varies is the ordering of the five steps above.
- If the copier stalls completely, what happens on the source cluster?Writers and readers on the source carry on unaffected - the hop is outside their path. What grows is the uncopied tail: records piling up on the source that the target does not have. The one indirect effect is that a copier reading hard after a stall competes for the source's read bandwidth with ordinary readers.
- Does adding more copier workers eventually make copy lag between the two clusters zero?No. More workers raise the hop's throughput, which stops the gap growing and shrinks it after a backlog, but the floor is set by the read cadence, the inter-site round trip and the target's own write. The source has already acknowledged the writer before any of that begins.
- Is a copy on a second cluster the same thing as keeping more copies of a partition inside one cluster?No. In-cluster copies are part of the cluster's own write path and a write can be made to wait for them. A cross-cluster copy is a separate cluster fed by an outside process after the fact. They protect against different failures and carry completely different timing.
A courier photographing the pages of a ledger that the clerk is still writing in. However fast the courier works, the copy in the other office is always a few lines short, and the clerk never pauses to wait for him.
saying these in an interview costs you the question
- Assumes the writer's acknowledgement on the source means the record is already on both clusters
- Says the copy is synchronous because the two clusters are linked to each other
- Claims copy lag between the two clusters can be tuned to zero with a faster link
- Treats any non-zero copy lag as a fault to be driven to zero rather than a measured property
- Confuses the copy on the target cluster with an extra in-cluster copy of a partition
- Believes a stalled copier blocks writers on the source cluster