skip to content

Why can an object store's asynchronous copy in a second region be missing objects that ingest already wrote successfully?

level: seniorimportance: should knowfreq 45%

answer

  1. acknowledged first, copied later
  2. the queue sets the distance
  3. a burst outruns the drain rate
  4. a trailing view, not a snapshot
  5. alert on oldest pending age

basics

~20 s

Because the write is acknowledged once the near copies are durable, and the far copy is applied afterwards by a background worker. That copy trails by a variable lag, so the newest objects have not reached it yet.

solid answer

~50 s

The acknowledgement and the far copy happen at different times. A write returns success once the copies inside the near scope exist, and the far-region copy is then queued and applied by a background process. How far behind that process runs is not a constant: it grows with the write rate, with object size, with the throughput available between the regions, and with any retries, so a burst of ingest can push the far copy from seconds behind to hours behind. The far copy is therefore a trailing view, not a snapshot - objects arrive in whatever order the workers drain them, so it can be a mixture of moments rather than a clean instant. Two practical consequences: reading the far copy can legitimately return not-found for an object you know was acknowledged, and the lag is something you should be measuring continuously rather than discovering during an incident.

code

pseudocode · 15 lines
pseudocode
on write(objectKey, bytes):
    write copies in the near scope
    wait until all near copies are durable
    acknowledge success to the client      // durability applies from here
    enqueue replicationTask(objectKey, target = secondRegion)

replication worker:
    for each task drained from the queue:
        copy objectKey to secondRegion
        mark objectKey as replicated       // may be minutes or hours later

read(objectKey) from secondRegion:
    if objectKey is not marked replicated:
        return notFound                    // acknowledged, not yet copied
    return bytes

go deeper

for a junior

Know the order of events: the write is acknowledged when the near copies exist, and the far-region copy happens afterwards in the background.

for a middle

Explain what moves the lag - write rate, object size and count, available throughput, retries - and why a burst can push it from seconds to hours.

for a senior

Show the operational view: measure the age of the oldest pending object, alert on the trend, and explain why not-found on the far copy is not a loss investigation.

for a principal

Own the trade explicitly. Say how many minutes of writes the far copy puts at risk under normal load and under the worst observed burst, and make that number a stated property of the platform rather than a discovery.

## Where the acknowledgement happens A write to an object store returns success at a specific moment: when the copies in the near scope are durable. From the caller's point of view the object is stored from that instant, and the durability promise applies to it from that instant. A copy in a second region is not part of that moment. It is created afterwards, by a background process that picks the object up and transfers it. This is a deliberate design: the distance between regions makes a synchronous far copy costly on every single write, so the product trades a strict guarantee for latency, and hands you a copy that is behind by a variable amount instead. ## What sets the lag The gap is not a published constant. It is the output of a queue, and it moves with: - **Write rate.** A landing zone that takes a steady stream all day and then a large batch drop will have two completely different lags in the same day. - **Object size and count.** Many small objects cost per-object work; a few huge objects cost throughput. Both can be the bottleneck, and they bottleneck differently. - **Available throughput between the regions.** Replication shares that path with everything else you move. - **Retries and failures.** A transient error puts an object back in the queue behind everything that arrived since. - **Anything that pauses the worker.** Once the queue is behind, it only catches up when the inbound rate drops below the drain rate, so a backlog built in twenty minutes can take hours to clear. ## The far copy is a trailing view, not a snapshot This is the part that surprises people the first time. The workers drain the queue in an order that suits throughput, not an order that reproduces your write history, so at any instant the far region can hold object A from ten minutes ago and object B from ten seconds ago. That has consequences beyond missing objects: - A set of objects written together as one logical batch can be **partly** present. - A reader that expects to see a complete hour will see a ragged edge at the end of it. - A deletion replicates through the same pipeline, so the far copy can still hold an object the near store no longer has. If a downstream consumer needs an all-or-nothing view, that has to be built on top - for example by writing a marker object only after every object in the batch is confirmed present, and having the consumer key off the marker rather than off the presence of individual objects. ## What you can observe The lag is measurable, and measuring it is the difference between an engineer who knows and one who guesses: 1. **Per-object replication state.** Whether a given object has been copied yet is normally exposed, which is what lets you explain a specific not-found rather than argue about it. 2. **Backlog size.** The number of objects, or bytes, still pending. 3. **Age of the oldest pending object.** This is the honest headline number, because it is the actual distance in time between the two copies. A backlog of ten objects whose oldest is four hours old is a worse situation than a backlog of ten thousand whose oldest is a minute old. Alert on the age rather than the count, and make the alert fire on a trend, since a backlog that is growing is a different problem from one that is large but draining. ## What it means for reading the far copy Two honest statements to make in an interview: - **Not-found on the far copy is not evidence of loss.** The acknowledged object exists in the near scope; it simply has not been copied yet. Treating that as data loss starts the wrong investigation. - **The copy protects against losing the region, at the price of the newest writes.** That is the trade you bought, and you should know roughly how many minutes of writes it puts at risk under normal load and under your worst observed burst - both numbers, because the second one is the one that matters when you need the copy. And one thing the far copy is not: it converges on whatever the near store currently holds, so it carries a bad write or a deletion across as faithfully as a good one. Its value is against the loss of the region, not against the loss of a day's work.

  • A nightly job reads the far copy and finds an hour of events only half present. What went wrong?
    Nothing failed - the queue was still draining. Replication applies objects in whatever order suits throughput, so the edge of a recent window is ragged by construction. If the job needs a complete window, it should wait on an explicit completeness signal, such as a marker written only once every object in the batch is confirmed copied.
  • Which signal tells you best how far behind the far copy is?
    The age of the oldest object still waiting to be copied. A count of pending objects says how much work is left but not how much time separates the two copies, and those diverge badly - a small backlog of stuck objects is a larger gap than a huge backlog that is draining fast.
  • Does a bad write reach the far copy as well?
    Yes, and a deletion does too. The far copy converges on whatever the near store currently holds, so it carries damage across faithfully after the same lag. Its protection is against losing the region, not against a mistake by whoever was writing.
  • Can you force the far copy to be current before a planned cutover?
    You can stop writing and wait for the backlog to drain to zero, which is the only reliable way to make the two copies agree. That requires a quiet period you have to plan for, and the time it takes is set by the backlog at the moment you stop, not by an average lag figure.

saying these in an interview costs you the question

  • Assumes the far copy is current the moment a write is acknowledged
  • Reads not-found on the far copy as evidence of data loss
  • Treats the replication lag as a fixed published number
  • Expects the far copy to be a clean point-in-time snapshot
  • Thinks a backlog clears as fast as it built up