After an unplanned failover on an asynchronously replicated relational cluster, how much committed data can be lost, how do you measure that exposure before an incident, and what levers reduce it?
answer
- RPO = committed but not shipped
- send vs receive vs apply lag
- measure bytes and seconds, p99
- promote furthest; cap lag for auto-promotion
- sync commit: latency + standby-down policy
basics
~20 sYou can lose everything the primary committed but had not yet shipped to the promoted replica. Measure it as replication lag in bytes and seconds at the p99, per replica, continuously. Levers: reduce lag, promote the furthest replica, refuse promotion past a lag threshold, or make commits synchronous.
solid answer
~60 sThe loss window is the gap between local commit durability on the primary and receipt on the replica that gets promoted. Under asynchronous replication that gap is real by design: the client is told committed before any remote copy exists. Measure it two ways, because they answer different questions. **Bytes** (WAL/LSN or binlog position difference) tell you how much data is at risk; **seconds** tell you how far behind in time. Track send, receive, and apply lag separately: a replica may have received everything but be slow to apply, which affects promotion time rather than loss. Alert on the p99, not the mean, and remember lag spikes exactly when you are most likely to fail over: bulk loads, purge storms, long transactions, network saturation. Levers, in order of cost: keep lag small (faster apply, avoid huge single transactions); always promote the furthest-ahead candidate; configure the manager to refuse candidates beyond a maximum lag; and make some commits synchronous so acknowledgement waits for a remote copy, which bounds RPO to near zero at the cost of write latency and a policy for what happens when the standby is unavailable.
go deeper
Know that asynchronous replication means unshipped commits are lost on a crash failover, and that lag is the thing to watch.
Distinguish send, receive and apply lag, and describe measuring lag in both bytes and seconds.
Talk in percentiles, tie lag spikes to specific workloads, cap automatic promotion by lag, and lay out the synchronous-commit tradeoff including the standby-down behaviour.
Turn it into a stated RPO per workload class with the latency budget it costs, decide where synchronous commit is worth it, and make the degrade-or-stall policy an explicit, signed-off choice.
## Where the exposure comes from A write on the primary becomes durable when its log record is flushed locally. Under asynchronous replication the client's COMMIT returns at that point. The log record is then streamed to replicas, received, and applied. Every stage adds delay, and a primary that dies mid-flight takes with it everything not yet received elsewhere. So the recovery point objective of a crash failover is not a configuration value; it is an observed property of your lag distribution at the worst moment. ## The three lags, and why the distinction matters - **Send/flush lag**: log written on the primary but not yet sent, or not yet flushed on the replica. This is the part that is truly lost on a crash failover. - **Receive lag**: data received by the replica but not yet made durable there. - **Apply/replay lag**: received and durable, but not yet replayed into the visible data. Not a loss risk, but it delays promotion (the node must finish applying) and makes read replicas serve stale data. A replica can be zero on send lag and minutes behind on apply lag, for example while replaying a large index build or blocked by a conflicting long-running query. Conflating the two produces both false alarms and false comfort. ## Measuring before the incident Practical practice: - Record lag in **bytes** and in **seconds** per replica, continuously, and keep the history. Bytes are the honest measure of data at risk; seconds are what stakeholders understand. - Alert on percentiles, since RPO is set by the bad moments, not the average. A cluster with 20 ms median lag and 40 s p99 has a 40 s RPO for practical purposes. - Correlate lag with the workloads that cause it: batch jobs, retention deletes, schema changes, backup windows, cross-site link saturation. - Use a heartbeat table or the engine's timestamp-based lag metric when idle periods make byte-based lag misleading (an idle primary shows zero lag whether or not the link works, so also alert on the stream being disconnected). - Rehearse: run a failover in a production-like clone under representative load and measure the actual gap between the last transaction the application saw and the last one present after promotion. ## Levers 1. **Keep lag small.** Fix the causes: split huge transactions, throttle bulk deletes, ensure replicas have IO headroom, avoid single-threaded apply bottlenecks (use parallel apply where available), give replication its own bandwidth. 2. **Promote the furthest candidate.** Automated managers compare positions; make sure your policy does too. Promoting a convenient node rather than the most advanced one wastes RPO. 3. **Bound promotion by lag.** Configure a maximum acceptable lag for automatic promotion; beyond it, escalate to a human rather than silently discarding minutes of transactions. 4. **Synchronous commit.** Have the primary wait for at least one replica to confirm receipt (or durability) before acknowledging. This bounds RPO to near zero for those transactions, at the cost of adding a network round trip to every commit and creating an availability question: if the synchronous standby is down, does the primary stall (safe, unavailable) or fall back to asynchronous (available, exposed)? That fallback choice is a business decision and should be explicit. 5. **Per-transaction granularity.** Many systems let you require synchronous commit only for the transactions that warrant it, for example payments, and leave the rest asynchronous. This targets latency cost where the value is. 6. **Retain the dead node's disk.** If the failed primary's storage survives, its unshipped transactions can sometimes be extracted forensically. Not a plan, but a reason not to wipe the node immediately. ## What to tell the business Express RPO as a number with conditions: under normal load we lose under a second; during the nightly batch, up to N seconds; with synchronous commit on payment transactions, none of those are lost. That framing survives an incident review far better than a claim of near-zero. ## Interview framing Name the window precisely (committed locally, not received remotely), separate send from apply lag, insist on percentile measurement in bytes and seconds, then give the ladder of levers ending with synchronous commit and its explicit availability tradeoff.
- You enable synchronous commit and the synchronous standby goes down. What should happen?You choose in advance between two behaviours: block writes until a synchronous standby is available, preserving the zero-loss guarantee at the cost of an outage, or automatically degrade to asynchronous, preserving availability while silently reopening the loss window. A common middle ground is to configure multiple candidate synchronous standbys so a single failure does not force the choice. What matters is stating that the default is a deliberate policy decision, not an implementation detail.
- Your monitoring shows zero replication lag but you are not confident. Why might that be misleading?On an idle or low-write primary, byte-based lag is zero regardless of whether replication is actually flowing, so a broken stream can look perfectly healthy. Similarly, a replica can show zero receive lag while being far behind on apply, which does not risk data but does delay promotion and serves stale reads. The fix is to monitor stream connectivity and a periodic heartbeat write in addition to positional lag, and to track apply lag separately.
saying these in an interview costs you the question
- Quoting a zero RPO for an asynchronously replicated cluster
- Monitoring only average lag instead of the tail
- Treating apply lag and send lag as the same number
- Assuming synchronous commit is free rather than a per-commit round trip
- Trusting zero-lag readings on an idle primary without a heartbeat or connectivity check