skip to content

Both sites accepted writes to a stream of the same name during an hour-long network split. Why can no merge afterwards recover the true order across them?

level: seniorimportance: should knowfreq 45%

answer

  1. two sequencers, not one
  2. nothing recorded the interleaving
  3. a merge invents a rule
  4. timestamps are two clocks
  5. route one entity to one site

basics

~20 s

There was never one order to recover: each site sequenced only its own writes, and nothing anywhere recorded how the two interleaved. Any merge has to invent a rule — write timestamps, site priority — and different rules produce different, equally unfounded answers.

solid answer

~40 s

A stream's order is produced by a single sequencer: the cluster that accepted the writes and decided what followed what. During the split there were two sequencers, so there are two real sequences and no artefact describing their interleaving — no counter advanced across both sites, and nothing at one site observed a write at the other. A merge afterwards therefore imposes an order rather than recovering one. Sorting on write timestamps sorts two clocks' opinions, and the skew can exceed the gap between the very writes you care about; picking a winning site is repeatable but arbitrary. The operable answer is routing: accept everything touching one entity at one site, so that entity has a single history, and leave only cross-entity order undefined.

go deeper

for a junior

Recall that each cluster numbers and orders only the writes it accepted itself. Two clusters accepting writes at the same time produce two separate sequences, and nothing recorded which write happened first across them.

for a middle

Explain why position numbers from two clusters are not comparable and why timestamps are two clocks rather than one shared sequence. Be able to say that a merge imposes an order instead of recovering one.

for a senior

Demonstrate the production consequence: two sites end up with different interleavings of the same records, so state folded from the stream differs by site and looks like a downstream bug. Then name the routing rule that avoids the whole question.

for a principal

Frame the merge rule as a business decision about which accepted writes may be silently dropped, and press on whether the organisation genuinely needs both sites accepting writes for the same entity or has asked for something else.

## What 'order' means before anything crosses a boundary A stream's order is whatever its platform promises locally, and the promises differ: some platforms split a stream into parts and guarantee order within a part, some guarantee it per key group, and some make no ordering promise at all between competing consumers. Whatever the promise is, it is produced by a **single sequencer** — the one cluster that accepted the writes and decided what followed what. That single sequencer is the only reason the local order is a fact rather than an opinion. ## What the split actually produced For that hour, each cluster was its own sequencer. Site A produced a sequence of A's writes; site B produced a sequence of B's writes. Both sequences are real, both are intact, and copying them around afterwards damages neither. What does not exist anywhere is the third thing people assume they can reconstruct: a record of **how the two interleaved**. No counter advanced across both sites. No artefact on either cluster observed a write at the other. The interleaving was never captured because nothing was in a position to capture it. ## Why a merge imposes an order rather than recovering one 1. **There is no shared sequence.** A record's position number belongs to the cluster that appended it — a receiving cluster assigns its own as the copier writes — so numbers drawn from two clusters are not comparable quantities. 2. **Timestamps are two clocks, not one.** Sorting on the write timestamp is the usual proposal, and what it sorts is the two clocks' opinions. Where the clocks disagree by more than the gap between two competing writes, the sort is simply wrong, and that gap is smallest exactly where the ordering matters most. 3. **Nothing recorded causality.** No write at one site carried evidence that its author had, or had not, seen a particular write at the other. Without that, 'concurrent' is the honest description of every cross-site pair. So a merge is a rule someone chooses, and the available rules disagree with each other: | Merge rule | What it gives you | What it quietly assumes | |---|---|---| | Sort on the write timestamp | one sequence, easy to build | the two sites' clocks agree more finely than the gap between competing writes | | One site wins, always or on ties | a deterministic, repeatable sequence | that site's hour genuinely outranks the other's | | Keep one value per entity, discard the rest | a single current state | losing an accepted write is acceptable here | | Keep both versions, hand them to the application | an honest result | every downstream can express what two concurrent versions mean | The lower half of that table is conflict resolution, which has its own literature and its own home. What belongs to an operator is the recognition that picking a row is a **decision about which accepted writes may be silently dropped**, not a wiring detail of the copy hops. ## What readers see afterwards Once the network recovers and both sequences have been copied in both directions, each site holds its own writes plus the copies, interleaved by whenever its copier delivered them. Two things follow. - The interleaving is **not identical at the two sites**, because two independent copiers ran at two different speeds against two different sources. A downstream that folds records into state can therefore reach a different result depending on which site it read from — a difference that survives every retry and looks exactly like a bug in the downstream. - Readers absorb **duplicates**: the echo, if the hops were never stamped with an origin, and the re-appends of a copier that resumed from an earlier point. Naming that as a consequence belongs here; how a consumer is built to tolerate it is a separate subject with a separate owner. ## The routing rule that makes it stop mattering The operable answer is almost never a cleverer merge. It is to arrange that everything touching one entity is accepted at **one** site. That entity then has a single sequencer and an intact history, and what remains undefined is the order of writes to *different* entities — something very little downstream state depends on. Pushed one step further, the same reasoning produces one stream per site with a single owning writer, read everywhere, instead of one shared name written in two places. ## Where platforms differ - What survives locally varies: order per part, order per key group, or no cross-consumer order at all. Saying 'order is preserved within a partition' as if it were universal is describing one design, not the class. - Whether a bad merge can be re-run varies: where a reader owns a rewindable position, the merged view can be rebuilt from the streams; where reading removes the record, there is nothing left to re-run against. - Renumbering varies, but the conclusion does not. Some platforms renumber positions across a hop and some preserve gaps, and neither creates a sequence that spans two clusters.

  • Does keeping all writes for one entity at a single site remove the problem?
    It removes the part that usually matters. If every write touching a given entity is accepted at one site, that entity has one sequencer and an intact history; what stays undefined is the relative order of writes to different entities, which most downstream state does not depend on. Choosing that routing rule is the real design decision behind any arrangement that accepts writes at two sites.
  • The two sites' clocks are synchronised to a few milliseconds. Is a timestamp merge good enough then?
    Only if competing writes at the two sites are separated by much more than the clock error, which is exactly what this arrangement is least able to promise. Synchronisation narrows the window; it does not create a shared sequence. The merged order is still a rule you chose rather than a record of what happened, and it will be wrong precisely on the close calls.
  • What does a reader at one site see once both sequences have been copied both ways?
    Its own site's writes plus the copies, interleaved by whenever its copier delivered them — and that interleaving need not match the other site's. A downstream folding records into state can therefore reach a different result per site. On top of that, readers absorb whatever duplicates the hops produced, from a missing origin stamp or a copier that resumed from an earlier point.

saying these in an interview costs you the question

  • Says sorting both sequences by record timestamp restores the true order
  • Assumes arrival order on the receiving cluster reflects real write order
  • Treats a globally unique record identifier as if it carried order
  • Claims the problem disappears once the network recovers
  • Says the two sequences merge cleanly because both are append-only
  • Insists a stricter in-cluster durability setting would have prevented it