Database Replication Explained: Replicas, Lag and Failover
8 min readBytePatterns
Database replication explained for system design interviews: primary and read replicas, replication lag, read-your-writes, sync vs async, and failover loss.
"Add read replicas" is a standard move in a system design interview, and it carries the same hidden cost as a cache: a second copy of the data that is sometimes behind. Replication gives you read capacity and a standby for when the primary dies. It does not give you write capacity, and it introduces replication lag, the gap that makes users see their own changes disappear. Interviewers probe which reads can tolerate that gap and what happens to committed data on failover.
The problem it solves
One database server has a ceiling, and for read-heavy products the reads hit it first. Leader-follower replication, also called primary-replica, keeps extra copies of the whole database. All writes go to the primary, which records each change in an ordered log. Replicas receive that log and apply the changes in the same order, so each one is a copy of the primary as of some recent moment. Reads can go to any copy.
That buys read throughput that grows with the number of replicas, reads closer to users in other regions, and a warm standby that can be promoted when the primary fails.
The intuition
Every change has a position in the primary's log, often called a log sequence number or LSN. A replica's state is fully described by one number: how far into the log it has applied.
Asynchronous replication commits on the primary and acknowledges the client immediately; replicas catch up in their own time. Writes are fast, but a replica can be behind by milliseconds under normal load and by much more under heavy load. A user who saves a change and is then routed to a lagging replica reads the old value. The fix most systems use is read-your-writes routing: remember the log position of the user's last write and only read from a replica that has applied at least that far, or from the primary.
Synchronous replication makes the commit wait for replicas to confirm. That removes the lag for the replicas that confirmed, at the price of write latency set by the slowest one waited for, and it blocks writes if that replica is down. Most setups compromise: wait for one replica, not all. The details matter: as of September 2026, PostgreSQL's synchronous_commit = on waits for the standby to make the change durable, and only remote_apply waits until it is visible to reads there; MySQL's semi-synchronous mode waits for a replica to receive the change, not to apply it.
Failover is where asynchronous lag becomes data loss. If the primary dies, a replica is promoted, and any commits it had not yet received are gone, even though clients were told they succeeded.
Watch it run
The animation opens on one writer, many readers, and a small delay. Every write is committed on the primary first, then streamed out to each replica, which applies it in its own time. Reads can be served from whichever copy is free: three times the read capacity. But writes still funnel through one machine, so write capacity stays at one. Under load the stream falls behind: the primary is on version 9 and replica 2 is not, 240 ms back. So a user can save a change and immediately read back the old value. Synchronous replication removes that, because the write waits for a replica to confirm, at 18 ms extra write latency. Then the primary dies, replica 1 is promoted, and it becomes the new truth. Fail over to a replica that was behind and the committed writes it never received are gone: three commits. The closing frame is the design decision: which reads are allowed to be stale, decided per query.
Database Replication
Step 1 of 12
Replication keeps extra copies in step. One writer, many readers, and a small delay.
The same interactive animation as the lesson — step through it with the controls.
The code
A toy model, not a real database: a primary that appends every write to an ordered log, and replicas that apply a prefix of it. A lagging replica serves the user's own write back as the old value:
class Primary:
def __init__(self):
self.log, self.data = [], {} # log entry: (lsn, key, value)
def write(self, key, value):
self.data[key] = value
self.log.append((len(self.log) + 1, key, value))
return len(self.log) # the write's log position
class Replica:
def __init__(self, name):
self.name, self.applied, self.data = name, 0, {}
def catch_up(self, primary, upto=None):
target = len(primary.log) if upto is None else min(upto, len(primary.log))
for lsn, key, value in primary.log[self.applied:target]:
self.data[key] = value # same order as the primary
self.applied = max(self.applied, target)
p = Primary()
r1, r2 = Replica("r1"), Replica("r2")
p.write("bio", "v1")
r1.catch_up(p)
r2.catch_up(p)
token = p.write("bio", "v2") # the user saves a new bio
r1.catch_up(p) # r1 keeps up, r2 lags behind
print(r1.data["bio"], r2.data["bio"], token, r2.applied) # v2 v1 2 1
Read-your-writes routing: remember the log position of the session's last write and only read from a copy that has applied at least that far. Failing over to a lagging replica loses whatever it never received:
def read(key, token, primary, replicas):
fresh = [r for r in replicas if r.applied >= token]
source = fresh[0] if fresh else None # any caught-up replica will do
return (source.data if source else primary.data).get(key), (source.name if source else "primary")
print(read("bio", token, p, [r2, r1])) # ('v2', 'r1')
print(read("bio", token, p, [r2])) # ('v2', 'primary')
for i in range(3, 13):
p.write("bio", f"v{i}")
r1.catch_up(p, upto=9) # the primary dies at log position 12
survivor = max([r1, r2], key=lambda r: r.applied) # promote the most caught-up copy
print(survivor.name, len(p.log) - survivor.applied, survivor.data["bio"]) # r1 3 v9
Write latency under the three commit modes, with made-up replica acknowledgement times: asynchronous waits for nobody, semi-synchronous for the fastest replica, fully synchronous for the slowest. One slow replica makes every write slow:
local_commit, acks = 2, [18, 25, 240] # ms; the third replica is far away
print(local_commit, local_commit + min(acks), local_commit + max(acks)) # 2 20 242
Checked on 400 seeded random histories of writes, partial catch-ups and reads by 5 users: every replica must equal a brute-force replay of the primary's log up to its applied position, reads routed by token must never return anything older than the user's own last write, and reads sent to a random replica must be caught doing exactly that:
import random
def replay(log, upto):
"""Brute force: rebuild a copy from scratch by replaying the log prefix."""
state = {}
for lsn, key, value in log[:upto]:
state[key] = value
return state
random.seed(26)
ok, naive_stale = True, 0
for _ in range(400):
prim, reps = Primary(), [Replica(f"r{i}") for i in range(3)]
last = {} # user -> (token, value written)
for _ in range(60):
user, action = random.randrange(5), random.random()
if action < 0.4:
value = random.randrange(1000)
last[user] = (prim.write(f"u{user}", value), value)
elif action < 0.7:
r = random.choice(reps)
r.catch_up(prim, upto=r.applied + random.randint(0, 3))
elif user in last:
tok, value = last[user]
ok &= read(f"u{user}", tok, prim, reps)[0] == value
naive_stale += random.choice(reps).data.get(f"u{user}") != value
ok &= all(r.data == replay(prim.log, r.applied) for r in reps)
print(ok, naive_stale > 0) # True True
The complexity
- Reads: capacity grows roughly with the number of copies, for reads that tolerate lag.
- Writes: unchanged at best; every replica applies every write too.
- Lag: usually small, unbounded in theory; write bursts and slow replicas grow it.
- Synchronous commit: local commit plus the slowest acknowledgement you wait for.
Where it goes wrong
- Expecting replicas to scale writes. They do not; that is partitioning, covered in database sharding.
- Routing every read to a replica. A user's own recent writes belong on the primary or behind a token.
- Monotonic reads. Two reads from replicas at different lag can make data go back in time. Pin a session to one replica.
- Failover without a policy. Promoting a lagging replica discards acknowledged commits. Accept that, or wait for a replica on commit.
- Split brain. If the old primary comes back still accepting writes, two copies diverge. Failover must fence it off.
When it shows up in interviews
Any read-heavy design uses it: news feeds, product catalogues, URL shorteners. The follow-ups are predictable: why did the user not see their update, what happens when the primary dies, sync or async, and how replication relates to the CAP theorem. It pairs naturally with caching: both are second copies of the truth, and both force you to say how stale is acceptable.
How to say it in an interview
"Writes go to the primary, which streams its log to replicas; reads go to replicas. That scales reads, not writes, and adds lag. For the user's own data I track the log position of their last write and only read from a replica that has applied it, otherwise from the primary. Replication is asynchronous by default for write latency, with one synchronous replica if we cannot lose acknowledged commits on failover. On failover I promote the most caught-up replica and fence the old primary so it cannot accept writes."