Your primary database node dies and a replica is promoted, yet several transactions the application had already been told were committed are missing. How is that possible, and how would you configure the system so it cannot happen?
answer
- Local durable ≠ survives losing the node
- Hole = replication lag at failure
- Sync/semi-sync: replica acks before client
- Quorum k-of-n keeps availability
- Fence old primary; promote the most advanced
basics
~20 sWith asynchronous replication, the primary acknowledges a commit once it is durable locally; replicas receive it slightly later. If the primary dies before shipping those transactions and you promote a replica, they are lost. Preventing it means requiring at least one replica to acknowledge before the client is told committed.
solid answer
~60 sDurability is scoped to a failure domain. Asynchronous replication makes the commit durable **on the primary's storage**, then streams it onward; the replica lags by microseconds to seconds. If the primary is lost permanently — hardware death, terminated instance — and you promote the lagging replica, everything in that lag window is gone, even though clients were told 'committed'. This is expected behaviour, not a bug, and the lost amount is your recovery-point objective. To close it, move the durability boundary: **synchronous commit** requires at least one replica to acknowledge before the primary answers the client. Then losing the primary loses nothing that was acknowledged. The costs are real: every commit pays a round trip to the replica (sub-millisecond same-rack, single-digit milliseconds cross-zone, tens or more cross-region), and a stalled or dead replica can block commits unless you use a quorum — acknowledge when any k of n replicas confirm — which keeps both durability and availability as long as a quorum survives. Also decide what 'acknowledged' means on the replica: received into memory (fast, tiny window) versus made durable there (safer, slower). And fence the old primary so it cannot come back and accept writes.
code
text · 14 linesasync replication
client -> primary: COMMIT
primary: durable locally
primary -> client: OK <-- acknowledged here
primary -> replica: ship record (may never happen)
data loss on failover = replication lag
synchronous / quorum
client -> primary: COMMIT
primary: durable locally
primary -> replicas: ship record
replicas -> primary: ACK (k of n)
primary -> client: OK <-- acknowledged here
data loss on failover = none for acknowledged transactionsgo deeper
Understand that a commit acknowledged by the primary may not be on the replicas yet, so a failover can lose recent data.
Explain that the lost window equals replication lag, and that synchronous acknowledgement before answering the client removes it at the cost of latency.
Design the whole configuration: quorum, received-versus-flushed acknowledgement, fencing, promotion of the most advanced replica, and lag as a monitored signal.
Set the recovery-point target per dataset, place replicas in failure domains that match the risk, price the commit-latency budget against it, and require tested failover rather than documented failover.
## Why acknowledged transactions can vanish Single-node durability answers 'will this survive a crash of this machine?'. It says nothing about 'will this survive the *loss* of this machine?'. With asynchronous replication: 1. Client commits; primary makes the record durable locally; primary answers 'committed'. 2. The record is streamed to replicas whenever the shipping process gets to it. 3. Primary is destroyed. Replica is promoted at whatever position it had reached. 4. Every transaction between the replica's position and the primary's last commit is lost. The size of that hole is the replication lag at the moment of failure — often milliseconds, but seconds or minutes under write bursts, a saturated network, or a replica busy with long queries. Lag is therefore a recovery-point metric and belongs on a dashboard with an alert, not just in a debugging session. ## Making it impossible: move the boundary **Synchronous commit / semi-synchronous replication.** The primary does not answer the client until at least one replica confirms it has the transaction. If the primary then dies, the surviving replica has everything that was acknowledged, so promotion loses nothing acknowledged. Two knobs define the strength: - **How many replicas must confirm.** One is enough to survive losing the primary. A **quorum** (k of n) also survives losing a replica without blocking commits, which is why production systems prefer quorum over 'one specific replica'. - **What confirmation means.** *Received into memory* is fast and leaves a tiny window (simultaneous loss of primary and replica memory). *Written and flushed on the replica* is stronger and slower. Most systems offer both; pick per workload. ## The costs you must state - **Latency on every commit.** The replica round trip is added to the critical path. Same rack: a fraction of a millisecond. Cross availability zone: typically single-digit milliseconds. Cross region: tens to well over a hundred. Cross-region synchronous commit is a deliberate choice for a small set of data, not a default. - **Availability coupling.** With exactly one required replica, that replica becomes a dependency of every write: if it stalls, writes stall. Quorum configurations remove the single dependency; without quorum you need a documented degrade policy (fall back to async and accept the changed recovery point, with alerting) rather than an operator improvising during an incident. - **Throughput.** The added latency shrinks per-connection commit rate, partially recovered by concurrency. ## The other half: failover correctness Synchronous replication guarantees the data is somewhere; it does not by itself guarantee the failover *uses* it. - **Fencing.** The old primary must be prevented from accepting writes after promotion. Otherwise a network partition produces two primaries, both accepting writes — split brain — and reconciling divergent histories is far worse than a small data-loss gap. - **Promote the most advanced replica.** In a quorum system, promotion must pick a node that has all acknowledged transactions; picking a lagging node throws away the guarantee you paid for. - **Client behaviour at failover.** In-flight commits whose acknowledgement was lost are indeterminate; applications need idempotent retries, not blind ones. - **Test it.** An untested failover path is an assumption. Regular exercises measuring actual data loss and time-to-recovery are the only way to know the configuration behaves as documented. ## Choosing the level Frame it per dataset as a recovery-point decision: - **Zero acknowledged loss required** (ledgers, payments, identity): synchronous with quorum, replicas in separate failure domains, fencing tested. - **Seconds of loss tolerable** (most application data): asynchronous replication with lag monitored and alerted. - **Regenerable data** (caches, derived tables, telemetry): async, and possibly relaxed local commit too. ## Answering well Explain the mechanism (local durability + lag = a hole equal to lag), name the fix (synchronous/quorum acknowledgement before answering the client), price it honestly (round trip per commit, availability coupling unless quorum), and then — the part that marks experience — add fencing, correct-replica promotion, lag as a monitored signal, and tested failover.
- With synchronous replication to exactly one replica, what happens if that replica goes down?Commits block, because the primary cannot get the acknowledgement it requires — you have coupled write availability to that one node. The standard remedy is a quorum over several replicas so any one can be lost without stalling writes. If quorum is not available, you need an explicit, alerting degrade policy that falls back to asynchronous and records the changed recovery point, rather than an operator deciding under pressure.
- Why is fencing the old primary essential after a failover?Without it, a primary that was only partitioned rather than dead can keep accepting writes while the promoted node also accepts writes — split brain. You then have two divergent histories with no automatic merge, and reconciling them means choosing whose committed data to discard. Fencing or a lease/quorum mechanism ensures the old node cannot serve writes after promotion.
- Does requiring the replica to acknowledge in memory rather than after its own flush weaken the guarantee?Slightly. It protects against losing the primary, which is the dominant failure, but a simultaneous power loss on both nodes could lose transactions held only in replica memory. It is a common and reasonable middle point because it removes the replica's flush latency from every commit; whether it is acceptable depends on whether the nodes share a failure domain such as a rack or a power feed.
saying these in an interview costs you the question
- Believing that having replicas automatically means no acknowledged transaction can be lost
- Not knowing that the loss window equals replication lag at the moment of failure
- Proposing synchronous replication with no mention of added commit latency or availability coupling
- Ignoring fencing, and therefore split brain, when describing failover
- Assuming failover automatically promotes the most up-to-date replica