Two nodes of an in-memory store each accepted writes as primary during a network partition, so when it heals why is one side usually discarded rather than merged?
answer
- two keyspaces, same names
- fence before anything else
- no version history per key
- opaque bytes have no merge rule
- the losing side is copied over
basics
~20 sThe tier keeps the current value of each key and no history of how it got there, and where values are opaque bytes there is no merge rule. So the losing node discards its keyspace, copies the winner's, and its acknowledged writes vanish.
solid answer
~50 sHealing leaves two keyspaces claiming the same names. Merging them would need two things this tier generally does not have: a record of which version came first, and a rule for combining two values. A store keeps the current value, not a version history, and where the server only hands back opaque bytes it has no way to combine two of them meaningfully. So the standard resolution is blunt — the losing node is fenced, demoted, and then **discards its own keyspace wholesale** and takes a full copy from the new primary. Every write it accepted during the split disappears, and those callers were already told the writes succeeded. Some products in this class are built to accept writes in several places and converge them using value types with a defined merge rule, but that is a deliberate, different design, not the behaviour of an ordinary single-primary tier.
go deeper
Recall that after a split the two nodes hold different data, and that one of them has to lose. The writes on the losing side do not come back.
Explain the two missing ingredients for a merge: no per-key version history, and no way to combine two values the server never interprets. Then say what actually happens instead, which is a full copy from the winner.
Show the operational shape: fence first, then discard, and bound the loss by the later of the node standing down and the last caller rediscovering the address. Say plainly that those callers were told yes.
Frame it as a policy question. Decide which keyspaces may ever be exposed to a silent discard of acknowledged writes, what the business cost of one is, and whether a product that converges concurrent writes is worth its constraints for the state you actually hold.
## Two keyspaces, one set of names When the link comes back, both nodes are holding a complete keyspace. They agree on the keys nobody touched, and they disagree on every key either side wrote during the split. A key may have one value on one node and a different value on the other; a key may exist on one and not on the other; a counter may have advanced independently on both. Something has to produce one keyspace. There are only two shapes of answer: combine the two, or pick one and throw the other away. This tier almost always does the second, and a senior answer explains why the first is not really on the table. ## Fencing comes first Before any resolution can be safe, the returning node has to stop being writable. **Fencing** here means exactly that: the old primary must be prevented from applying any further writes as though it were still in charge. The usual mechanism is that the tier's configuration carries a generation that rises on every promotion, and a node that sees a newer generation than its own stands down. Where an intervening proxy resolves keys, the proxy simply stops routing to it; where a controller outside the data plane decides, the controller stops or isolates it. The general idea of rejecting a superseded writer is borrowed from the distributed-systems material on electing a replacement, and is not re-derived here. What matters for this tier is the timing. Callers holding the old primary's address keep writing to it until they rediscover the new one, so the band of writes that will be thrown away runs from the start of the split to **whichever is later**: the old node standing down, or the last caller completing address discovery. The rediscovery window is frequently the longer of the two and is the one teams forget to measure. ## Why the tier cannot merge Merging requires answering two questions about every conflicting key, and an ordinary in-memory store can answer neither. - **Which value came later, causally?** The store holds the current value, not a chain of versions, and it keeps no per-key record of who wrote what and after which other write. Wall-clock timestamps are not a substitute: they exist only if the application put them in the value, and they compare clocks on two machines that were partitioned from each other. - **What would combining them even mean?** Where the server treats values as **opaque bytes** it hands back unchanged, there is no operation that combines two of them. Where the server does understand structure, a merge is conceivable for some shapes, but it still needs the causal history above to know what to merge, and the semantics belong to the application rather than to the store. Hence the blunt resolution: **the losing node discards its own keyspace and takes a full copy of the winner's state.** It is not that merging was attempted and failed; it was never available. ## What that costs, and to whom | What the affected keys held | Effect of discarding that side | Typical response | |---|---|---| | Derived state another system can rebuild | A period of extra work and slower responses while it refills | Accept it; size the refill load on the system of record | | Sole-copy ephemeral state such as claims, quotas or deduplication records | Silent, unrecoverable loss of things the caller was told were stored: a double execution, a reset quota, two holders of one claim | Refuse writes on the minority side instead, or keep the state elsewhere | The distinction is worth stating explicitly in an interview, because the same mechanism produces an inconvenience or an incident depending on the keys, and the mechanism itself cannot tell which it is. ## What makes this tier sharper than a durable engine On a durable engine a divergent branch usually still exists on disk after the fact, so the writes can at least be extracted and reconciled by hand later, however unpleasant that is. Here the discard is normally the end of it. The losing node's memory is overwritten by the incoming copy, and while some deployments do keep the node out until an operator has dumped its keyspace, and some products keep a write log on disk that survives, nothing feeds either back into the running primary. Reconciliation from a dump is application work: you would have to know what those writes meant and how they relate to what the winner recorded in the meantime, and the store cannot tell you. Say the consequence out loud, because it is the whole point: **the callers were told yes.** Acknowledgment on this tier is not a promise of durability, and a split brain is the case where that gap is not a few milliseconds but the entire length of the partition. ## The variation to keep in the answer Do not assert that stores of this class can never converge two writable copies. Some are built for precisely that, using value types whose merge rule is defined in advance so concurrent writes combine deterministically. That is a property you buy deliberately, it constrains what you may store and what guarantees you get back, and it is not what an ordinary single-primary deployment does when a partition heals. ## The band of lost writes is bounded by whichever is later: the old primary standing down, or the last caller finishing address discovery. The second is often the longer of the two ``` t0 network splits: hall A holds a majority, hall B holds the old primary t0+3s hall A promotes a replica; two nodes now accept writes for one keyspace t0+3s.. callers in hall B keep writing to the old primary and are told yes t0+90s link returns; the old primary sees a newer configuration generation t0+90s it stands down, discards its keyspace and copies the new primary's state t0+94s the last caller in hall B completes address discovery and stops writing to it after every write accepted in hall B from t0 to t0+94s is gone, and nobody tells those callers ```
- The losing node kept a write log on disk during the split. Can those writes be recovered?Not automatically. Nothing feeds a losing branch's log into a running primary, and on rejoin the node's state is replaced by a copy of the winner's. At best an operator keeps the node out of service and dumps the log for offline inspection. Turning that into recovered state is application work, because it requires knowing what each write meant relative to what the winner recorded during the same window.
- Why does address discovery matter to the size of the loss, if the old primary has already stood down?Because the two are not simultaneous. Between the moment the old node stops accepting writes and the moment the last caller stops addressing it, those callers see failures rather than silent acceptance, which is better — but before it stands down, every caller still pointed at it is having writes accepted that will be discarded. The band you lose ends at the later of the two events, so both windows need measuring.
- Would giving every value a timestamp let the tier merge instead of discard?No, not reliably. The timestamp would be written by the application, and it compares clocks on two machines that were by definition unable to talk to each other, so skew decides outcomes. It also picks one whole value per key rather than combining them, which silently drops the other side's change. Picking a winner by timestamp is a conflict-resolution strategy with well-known failure modes, not a property the store provides.
saying these in an interview costs you the question
- Assumes the tier merges the two keyspaces automatically once the link returns
- Says the losing writes can be replayed from the losing node's own log
- Thinks wall-clock timestamps on values make a safe merge possible
- Forgets the callers were already told the discarded writes succeeded
- Ignores fencing and lets the old primary keep taking writes while rejoining
- States flatly that no store of this class can converge concurrent writes