Database Replication Explained: Leaders, Lag, and Failover
Every design that survives a machine failure has replication underneath it, and almost every candidate who says "I'll add read replicas" stops one question short of the interesting part. Replicas are easy to draw. What they cost you is a stale read, a window of data loss on failover, and an automated promotion procedure that is the single most dangerous thing running in your infrastructure.
This post covers the mechanism — one ordered log, replayed — then the synchronous/asynchronous decision and its arithmetic, the three anomalies replication lag causes and how to fix each one, how failover actually goes wrong, and where replication stops helping.
The two reasons to replicate, and they pull apart
People say "replication" for two goals that want different configurations.
Read throughput. One primary cannot serve every read. Copies can. This is the reason most systems start replicating, and it is the weaker of the two: a cache in front of the database is usually cheaper per read served.
Surviving a machine. Disks fail, instances get retired, an availability zone loses power. A second live copy turns a data-loss event into a failover. This is the reason that actually justifies the operational cost.
They conflict because the first wants many cheap asynchronous copies spread wide, and the second wants at least one copy that is guaranteed current, which means paying a round trip on every commit. You can have both, but you have to say which replicas are which.
A third reason shows up in interviews and is worth naming separately: geographic read latency. A replica in Frankfurt serves European reads in 5 ms instead of 90 ms. It does nothing for European writes, which still cross the Atlantic to the leader.
| Goal | What it needs | What it costs |
|---|---|---|
| Read throughput | N async followers | Stale reads; write amplification on every copy |
| Surviving a machine | ≥1 sync or quorum copy | Commit latency = local write + a network round trip |
| Regional read latency | A follower per region | Cross-region bandwidth; lag proportional to distance |
| Point-in-time recovery | A delayed replica or backups | Storage; replication is not a backup (see below) |
The mechanism: one ordered log, replayed
Single-leader replication — what Postgres, MySQL, MongoDB and most managed services do by default — is simpler than the diagrams suggest. One node accepts writes. It writes them to a durable, strictly ordered log. Every other node pulls that log and replays it in the same order. That is the whole protocol.
What travels in the log matters more than it looks:
- Statement-based ships the SQL. Compact, and broken by anything non-deterministic —
NOW(),RAND(), anUPDATE ... LIMITwithout a deterministic order. MySQL defaulted to this for years and the bug reports are a genre. - Physical / WAL shipping ships the byte-level changes to disk pages. Exactly reproducible, which is why Postgres streaming replication is boring in the good way. The cost is that the replica must run the same major version and page layout.
- Logical / row-based ships "row with primary key 42 changed these columns to these values". Deterministic like physical, but decoupled from storage internals, so it replicates across versions and into things that are not the same database at all. This is what makes zero-downtime major upgrades and change data capture possible.
Say "logical replication" in an interview and you have opened the door to online schema migrations and CDC pipelines, which is usually where the follow-up wants to go.
Synchronous, asynchronous, and the quorum in between
This is the decision that determines everything else, and it comes down to one question: does the leader wait for a replica before telling the client the write succeeded?
The arithmetic is worth doing out loud. Suppose your leader commits locally in 1 ms and you take 5,000 writes per second.
Asynchronous, replica 200 ms behind. The leader dies. Writes committed in the last 200 ms exist nowhere else: 5,000 × 0.2 = about 1,000 acknowledged writes gone. If those were payments, you have a reconciliation problem, not an outage.
Synchronous to a replica in another AZ. Commit becomes roughly 1 ms + 2 ms round trip = 3 ms. Data loss on leader failure is zero. But a single connection doing sequential commits now tops out at ~330/s instead of ~1,000/s, so you need three times the write concurrency for the same throughput.
Synchronous across regions. Virginia to Ireland is about 75 ms round trip, and physics is not negotiable — that is roughly the speed of light in fibre over 5,500 km, twice. A sequential writer now manages about 13 commits per second. This is why nobody runs cross-region synchronous replication for a transactional workload, and why saying you would is a red flag.
Fully synchronous to all N replicas is worse than it looks for a different reason: availability. If every commit needs all three replicas, any one replica being slow or down blocks all writes. You have made three machines into a single point of failure.
The production answer is a quorum: wait for any k of n. Postgres spells it synchronous_standby_names = 'ANY 1 (r1, r2, r3)' — the leader waits for whichever replica answers first, so one slow node costs nothing and the write survives a single machine loss. MySQL's semi-synchronous mode does the same with rpl_semi_sync_source_wait_for_replica_count, with a caveat worth knowing: on timeout (default 10 seconds) it silently falls back to asynchronous rather than refusing writes. Your "synchronous" cluster is asynchronous during exactly the incident you bought it for.
One more Postgres detail that catches people: synchronous_commit = on waits for the standby to flush the WAL to disk, not to apply it. The write is durable but a read on that standby can still miss it. remote_apply closes that gap and costs more. If you promise read-your-writes from a synchronous replica, that setting is the promise.
Replication lag and the three anomalies
Async replication is eventually consistent, and "eventually" produces three distinct bugs. Naming all three is a strong interview signal, because most candidates know only the first.
Read-your-writes. A user updates their profile, the write goes to the leader, the page reloads, the read goes to a follower that has not applied it yet, and their change appears to have vanished. They hit save again. Now you have two writes and a support ticket.
Monotonic reads. Two successive reads hit different followers — the first at log position 108, the second at 103. Time appears to run backwards: a comment that existed disappears. The fix is to pin a given user to a given replica (hash the user ID to a replica, not round-robin) so their reads never move to a less-advanced cursor.
Consistent prefix reads. Causally related writes arrive out of order because they were replicated through different partitions. The classic shape is a reply appearing before the message it answers. Within one replicated log this cannot happen — the order is total — which is exactly why it starts happening once you shard and have one log per shard.
Notice that all three are artefacts of where you read, not of corruption. A lagging follower is never inconsistent with itself; it is a correct snapshot of an earlier moment. That framing is the one to use in an interview, and it maps directly onto the availability side of the CAP theorem: you chose to serve a possibly-stale answer rather than no answer.
Fixing read-your-writes, concretely
The lazy fix is "route all reads to the leader for 30 seconds after a write". It works, it is one line, and it throws away the replicas you paid for during the busiest part of a user's session. The precise fix is to carry the write's log position with the user and only read from a replica that has passed it.
# Write path: capture the leader's log position and hand it back to the caller.
def place_order(leader, user_id, items) -> tuple[int, str]:
with leader.transaction():
order_id = insert_order(leader, user_id, items)
token = leader.scalar("SELECT pg_current_wal_lsn()")
return order_id, token # store in the session, a cookie, or a header
# Read path: pick a replica that has replayed past the token, else use the leader.
def list_orders(pool, user_id, token):
if token is None: # user hasn't written recently
return select_orders(pool.any_replica(), user_id)
for replica in pool.replicas_shuffled():
# Compare as pg_lsn in the database — LSN text does not sort correctly
# as a string ('FF/0' would beat '100/0').
if replica.scalar("SELECT pg_last_wal_replay_lsn() >= %s::pg_lsn", token):
return select_orders(replica, user_id)
return select_orders(pool.leader, user_id) # fall back rather than serve stale
Three things make this work in practice. The token is per user, so one user's write never forces everyone to the leader. The fallback is to the leader, not to a stale replica — degrade latency, not correctness. And in a real deployment you do not issue that pg_last_wal_replay_lsn() query per request: a background poller refreshes each replica's position every 50 ms and the router compares against the cached value, which costs nothing and is at most 50 ms conservative.
MySQL has the same shape with GTIDs and WAIT_FOR_EXECUTED_GTID_SET(). DynamoDB collapses it into a per-request flag, ConsistentRead=true, which routes to the leader replica and costs double the read capacity units — the same trade, priced explicitly.
Measuring lag honestly
Seconds_Behind_Source in MySQL is not a measure of lag. It is the difference between the timestamp of the event currently being applied and the replica's clock — so an idle leader shows 0 even if the replica is minutes behind on unshipped events, and a broken replication thread shows NULL. It also goes wrong under chained replication.
The measurement that works is a heartbeat: a single row on the leader updated with the current timestamp every second, and the replica reading that row and subtracting. pt-heartbeat does exactly this and nothing else. The number it gives you is end-to-end time, which is the number that matters, because it directly answers "how much would I lose right now?" — multiply it by your write rate.
Alert on the heartbeat, page on it above a threshold you derived from that multiplication, and put the lag on the same dashboard as the write rate. Lag almost never grows because the network is slow; it grows because the replica's single apply thread cannot keep up with a leader that parallelised the writes across many connections. Postgres and modern MySQL both support parallel apply, and turning it on is usually the fix.
Failover is the hard part
Adding a replica is a config change. Promoting one is a distributed consensus problem you are solving under time pressure, and every part of it can go wrong.
A failover has four steps: detect that the leader is gone, choose a successor, redirect clients, and prevent the old leader from acting.
Detection is a timeout, and the timeout is a genuine dilemma. Short, and a GC pause or a brief network blip triggers an unnecessary promotion. Long, and you are down for the duration. Ten to thirty seconds is typical for managed Postgres and MySQL; consensus systems like etcd run at a default 1,000 ms election timeout with 100 ms heartbeats because their promotion is cheap and safe.
The step that produces incidents is the fourth.
Three failure modes to have ready:
Split-brain. Two nodes both believe they are leader and both accept writes. The defences are a quorum requirement for promotion (a minority can never elect) and fencing of the deposed leader — STONITH, a storage lease that expires, or a proxy layer like ProxySQL or PgBouncer that simply stops routing to it. Fencing tokens generalise this: every write carries a monotonically increasing epoch, and the storage layer rejects anything from an old epoch.
Lost writes. With asynchronous replication, promoting a follower discards everything the old leader committed but had not shipped. GitHub's October 2018 incident is the canonical worked example: a 43-second network partition triggered an automated failover, and reconciling the writes that had diverged left the site degraded for over 24 hours. The failover took 43 seconds of trigger and a day of consequences. That asymmetry is the thing to internalise.
The thundering herd after promotion. A new leader starts with a cold buffer pool and every connection in the fleet reconnecting at once. Plan for the first minute after failover to be slower than the outage it replaced.
Which is why the honest answer to "how do you handle failover?" is not "an automated tool". It is: quorum-based promotion, fencing, an explicit statement of how many writes you are willing to lose, and a runbook someone has actually rehearsed.
Multi-leader and leaderless, briefly
Multi-leader puts a writable leader in each region. Local write latency, survives a whole region, and it buys you write conflicts: two users edit the same row in Tokyo and Dublin within the replication window, and both writes are accepted. Something has to resolve them — last-write-wins (which silently discards data and depends on clock skew), version vectors, or CRDTs that merge by construction. Worth proposing for collaborative editing or offline-first mobile sync; a bad default for anything with an invariant, such as a balance that must not go negative.
Leaderless (Dynamo, Cassandra, Riak) removes the leader entirely: clients write to W replicas and read from R of N, and as long as W + R > N the read set overlaps the write set, so at least one replica returns the current value. With N=3, W=2, R=2, you survive one node down for both reads and writes. The price is that you do the repair work — read repair, anti-entropy, hinted handoff — and that quorums give you overlap, not linearizability. A read concurrent with a write can still return either value.
An interviewer who hears you distinguish "quorum overlap" from "strong consistency" will believe you have used these systems.
What replication does not solve
The write ceiling. Every replica applies every write. Say a machine handles 30,000 operations per second and you take 5,000 writes/s — each replica spends 5,000 on the write stream and has 25,000 left for reads, so four replicas give you 100,000 reads/s. Now push writes to 25,000/s: each replica has 5,000 ops of read capacity left, and the fifth replica you add makes the situation slightly worse, not better. Replication divides reads and multiplies writes. Moving the write ceiling is what sharding is for.
Backups. DELETE FROM orders with a bad WHERE clause replicates in single-digit milliseconds to every copy you own. Replication protects against a machine failing, never against a human or an application bug. A delayed replica (MySQL SOURCE_DELAY = 3600) is a cheap hedge that gives you an hour to notice, but it is not point-in-time recovery and it is not an off-site immutable backup.
Consistency. More copies means more places to read a stale value. Replication creates the anomalies above; it does not resolve them. That work is in your routing layer.
A bad query. A sequential scan on an unindexed column is just as slow on a replica. Replication multiplies your capacity to run the bad plan.
Write latency anywhere but one place. A European replica serves European reads fast and European writes exactly as slowly as before, because they still cross to the leader.
The alternatives worth naming
Consensus replication. Instead of a leader plus followers plus an external failover tool, run Raft or Paxos over the replica set: CockroachDB, TiKV, etcd, MongoDB's replica sets. A write commits when a majority durably has it, so an acknowledged write is never lost and promotion is part of the protocol rather than a script you hope works. The cost is a majority round trip on every write and a minimum of three nodes. If a question emphasises "we cannot lose an acknowledged write", this is the answer, not semi-sync.
Shared-storage replication. Aurora, Neon and AlloyDB replicate the storage layer instead of the query layer. Aurora writes each page to six copies across three AZs with a 4-of-6 write quorum and 3-of-6 for reads; replicas read the same storage rather than replaying a log, so lag is typically tens of milliseconds and adding a replica adds no write amplification at all. It removes most of the pain in this post, at the price of being a specific managed product.
Change data capture. If the reason you want a replica is to feed a search index, a cache or an analytics warehouse, do not add a database replica — stream the logical log with Debezium into the thing that actually needs it. Same log, cheaper consumer.
A cache. If the goal is purely read throughput, compare honestly against caching: a Redis node serves an order of magnitude more reads per pound than a database replica and costs you the same staleness you were already accepting. Replicas earn their place when reads must be fresh-ish, transactional, and expressible as arbitrary SQL.
Answering this in an interview
The sequence that lands:
- Separate the goals. "Are we replicating for read throughput, for surviving a machine, or for regional latency?" These want different topologies and the question usually does not say which.
- State the topology and the log. Single leader, followers replaying one ordered stream, logical or physical.
- Make the sync decision explicit and price it. "Async, so commits stay at 1 ms and we accept losing up to the lag window — at 5,000 writes/s and 200 ms of lag that's about 1,000 writes. If that's unacceptable for payments, quorum-sync to one of three replicas and pay 2 ms."
- Volunteer the lag anomaly before you are asked. Read-your-writes, and the LSN-token routing that fixes it properly.
- Treat failover as a design element, not a checkbox. Quorum promotion, fencing, split-brain, and what happens to the writes the old leader never shipped.
- Say where replication stops. It never raises the write ceiling and it is not a backup.
Step 3 is the one that separates candidates. Anyone can say "add read replicas". Quoting the loss window in writes rather than in milliseconds shows you have run this in production, or at least thought like someone who has.
If you want to hear how that sounds under pressure — including the follow-up about what happens to in-flight writes during promotion, which is where this question always goes — you can run a full system design mock by voice on Whitepad.
Practice this out loud
Reading is the easy part. Sit across from a senior AI interviewer that talks, watches your whiteboard, and scores you like the real thing — your first mock is free.
Start a free mock →