Replication
Keep copies of the same data on several machines. Copies survive a failed machine and serve more reads, but each copy is a little late.
- 1One leader takes writes and streams its log. Standbys replay it in order, a little late.
- 2Asynchronous replicas can lose the last commits on failover. A synchronous standby trades commit latency for zero loss.
- 3Replica reads can be stale. Fix read-your-writes with a token or a leader read after a write.
- 4Quorums need R + W > N, and a sloppy quorum breaks that rule.
- Replication copies every write from the leader to its replicas.
- Copies give failover, read capacity and copies near users.
- The cost: a copy lags, and a failover can lose what it never received.
A replica is always a little behind the leader. A read there can miss a change the same user just made.
| method | on leader failure | status |
|---|---|---|
| One leader, asynchronous standbys | The last commits can be lost. | Small loss accepted |
| One leader, quorum commit: ANY 1 of 2 standbys | No commit lost if a synchronous standby survives. | Approved |
| One leader, every standby synchronous | No loss, but one slow standby stops all writes. | Not approved |
| Multi-leader, one leader per region | Writes go on; concurrent edits conflict. | Multi-region writes |
| Leaderless quorum (N, R, W) | No failover step at all. | R + W > N |
| Replay SQL statements on each copy | now() and random() differ per copy. | Not approved |
Postgres streams its write-ahead log: the replica replays the exact bytes, not the statements.
I start with one leader and asynchronous standbys, and I add a synchronous standby when losing the last commits on failover is not acceptable.
| level | leader restart | standby OS crash | standby read | localhost |
|---|---|---|---|---|
| off | Loses 600 ms | May lose | Stale | 0.08 ms |
| local | Kept | May lose | Stale | 0.14 ms |
| remote_write | Kept | May lose | Stale | 0.35 ms |
| on | Kept | Kept | Stale | 0.36 ms |
| remote_apply | Kept | Kept | Fresh | 0.35 ms |
Guarantees from the Postgres 16 documentation, with a synchronous standby set. off can lose up to 3 × wal_writer_delay: 600 ms at the default 200 ms. Latency: median of 500 commits, primary and standby on one laptop; across regions add one network round trip.
I choose the durability level per transaction. Payments wait for a synchronous standby to flush; logs only wait for the local disk.
Step 1: Write
- Send every write to the leader.
- The commit returns at the synchronous_commit level: local disk, or a standby too.
If it fails
The leader is down. Writes fail until failover promotes a standby (step 5).
Writes go to the leader. Reads go to standbys, except a read that must see the user's own write, which waits for a token or goes to the leader.
| tool | capability | what it gives this design | also used for |
|---|---|---|---|
| Postgres | Streaming replication of the WAL | Each standby is an exact copy, replayed in commit order. | Point-in-time recovery from archived WAL |
| Postgres | Hot standby | Read-only queries on a standby while it replays. | Reports off the primary |
| Postgres | synchronous_standby_names: FIRST k or ANY k | A commit waits for k standbys. ANY 1 of 2 survives one slow standby. | Zero-loss failover |
| Postgres | synchronous_commit, settable per transaction | Durability per write: payments wait for a standby, logs do not. | Bulk loads with off |
| Postgres | pg_current_wal_insert_lsn(), pg_last_wal_replay_lsn() | A version token: the replica knows if it has a given commit. | Causal reads across services |
| Postgres | pg_stat_replication: write_lag, flush_lag, replay_lag | Lag per standby. The router skips a standby that lags too much. | Alerts |
| Postgres | recovery_min_apply_delay | A standby that stays behind on purpose. The lab used it to make lag visible. | Undo a bad DELETE |
| Postgres | Replication slots | Cap the size The leader keeps WAL until the standby has it. A dead standby fills the disk. | Logical decoding |
| Postgres | pg_promote() | Turns a standby into the leader. | Planned switchover |
| Redis | Asynchronous replicas; WAIT numreplicas timeout | Limit WAIT counts acknowledgements but does not make Redis strongly consistent. | Read scaling for caches |
| Cassandra | Consistency level per query: ONE, QUORUM, LOCAL_QUORUM, ALL | R and W chosen per read and write. Hinted handoff and read repair heal copies. | Multi-region writes |
| DynamoDB | ConsistentRead on a read | Choose a strongly consistent read; the default read can be stale. | Serverless key-value |
| Patroni · etcd | A leader key with a time to live | One leader at a time. The holder renews it; a new leader is promoted when it expires. | Any single-leader failover |
Postgres gives me ordered WAL streaming, per-transaction durability, and LSN functions for read-your-writes. Leaderless stores give me per-query consistency levels instead.
write(value, W):
send value to all N replicas
wait for W acks1 // the write set
read(R):
ask all N replicas
wait for R replies2 // the read set
RETURN the reply with the newest version3
R + W > N4 => every read set meets every write set- 1The write succeeds after W replicas store it. The others get it later, or never if they are down.
- 2Any R replicas. The read cannot know which ones have the write.
- 3Each value carries a version. One fresh reply is enough.
- 4R + W > N means the sets must overlap: there are only N replicas to choose from.
Tested source Go: every case, enumerated
// Outcome is what a write with W acknowledgements, then a read with R replies, can return,
// over every choice of which W nodes acknowledge and which R replicas reply.
type Outcome struct {
WriteOK bool // W nodes were reachable
ReadOK bool // R replicas were reachable
Cases int // pairs (write set, read set)
Stale int // pairs where no reply carries the write: the read returns the old value
}
// Enumerate tries every write set and every read set. A read sees the write when one of its R
// replies comes from a node in the write set; it returns the newest version among the replies.
func Enumerate(n, r, w int, p Pattern) Outcome {
wt, rt := targets(n, p)
o := Outcome{WriteOK: len(wt) >= w, ReadOK: len(rt) >= r}
if !o.WriteOK || !o.ReadOK {
return o
}
for _, acks := range subsets(wt, w) {
for _, reads := range subsets(rt, r) {
o.Cases++
if !overlaps(acks, reads) {
o.Stale++
}
}
}
return o
}
A write waits for W acknowledgements and a read for R replies. With R plus W greater than N, the two sets always share a replica, and the read takes the newest version.
| N, W, R | writes survive | reads survive | reads fresh |
|---|---|---|---|
| 3, 2, 2 | 1 down | 1 down | Approved |
| 3, 3, 1 | 0 down | 2 down | Approved |
| 3, 1, 3 | 2 down | 0 down | Approved |
| 3, 1, 1 | 2 down | 2 down | Stale in 6 of 9 |
| 5, 3, 3 | 2 down | 2 down | Approved |
Writes survive N − W failed replicas; reads survive N − R.
N = 3 with W = 2 and R = 2 survives one replica down for both reads and writes, and every read is fresh.
All replicas are up. The read arrives before the replicas outside the write set apply it.
Over all 9 choices of write set and read set, the read returns the old value in 0 (0%).
Every N from 1 to 5, every R and W, four patterns: 220 settings. Each enumerates all write sets and read sets.
I check two things: can each operation still gather enough replicas, and does every read set meet every write set.
| anomaly | the user sees | fix | cost |
|---|---|---|---|
| Read your writes | Their own change is missing after a reload. | Read the leader after a write; version token; remote_apply | Leader load; a short wait; slower commits |
| Monotonic reads | A comment appears, then disappears. | Sticky replica per user; token of the newest position seen | Uneven replica load; a short wait |
| Consistent prefix | An answer before its question. | Write related data to one partition, so one log orders it | The partition key choice |
Replica lag breaks three promises: users see their own writes, data never goes backwards, and an answer never appears before its question.
The user renames their profile, then reloads it 10 ms later. Each read goes to any replica, with no check.
Same case on a real Postgres standby that applies each commit 300 ms late: a read right after the write returned 'Ada'. A read that first waited for the token returned 'Grace'.
For read-your-writes I give the client its commit position. The replica waits until it has replayed that far, or the router reads the leader.
write(name): // leader
UPDATE profiles SET name = ...
token = WAL position // after the commit1
RETURN token
read(token): // replica
WHILE replay position < token2:
wait 2 ms // until a deadline3
SELECT name ...- 1Read the position in a second statement. A position read inside the transaction comes before the commit record.
- 2pg_last_wal_replay_lsn() on the replica, compared as pg_lsn.
- 3A deadline caps the wait when a replica is far behind. Then read the leader.
Tested source Go: write, then read after the token · SQL: token and replay check
// WriteName changes a name on the leader and returns a version token: the WAL position just past
// the commit. It reads the position after the commit, in a second statement; a position read
// inside the transaction comes before the commit record.
func WriteName(ctx context.Context, leader *pgx.Conn, id int, name string) (string, error) {
if _, err := leader.Exec(ctx, token["write_name"], id, name); err != nil {
return "", fmt.Errorf("write name %d: %w", id, err)
}
var lsn string
if err := leader.QueryRow(ctx, token["token"]).Scan(&lsn); err != nil {
return "", fmt.Errorf("read commit position: %w", err)
}
return lsn, nil
}
// ReadNameAfter reads a name from a replica once the replica has replayed the WAL up to the token,
// so the caller sees its own write. It gives up at the context deadline; the caller can then read
// from the leader.
func ReadNameAfter(ctx context.Context, replica *pgx.Conn, id int, lsn string) (string, error) {
for {
var caught bool
if err := replica.QueryRow(ctx, token["replayed"], lsn).Scan(&caught); err != nil {
return "", fmt.Errorf("check replay position: %w", err)
}
if caught {
break
}
select {
case <-ctx.Done():
return "", fmt.Errorf("replica has not replayed %s: %w", lsn, ctx.Err())
case <-time.After(2 * time.Millisecond):
}
}
return ReadName(ctx, replica, id)
}
UPDATE profiles SET name = $2 WHERE id = $1;
-- After the commit: the WAL position just past the commit record is the version token.
SELECT pg_current_wal_insert_lsn()::text;
-- On the replica: has it replayed the WAL up to the token?
SELECT pg_last_wal_replay_lsn() >= $1::pg_lsn;- right after
- 'Ada': the standby is 300 ms behind.
- with token
- 'Grace', after a wait of about 300 ms.
- remote_apply
- 'Hopper': the commit itself waited.
The commit position is a version token. A replica that has replayed past it has the write.
- Detect: the leader stops renewing its key. Patroni gives the key a 30 s time to live by default.
- Choose: the standby with the most replayed WAL loses the least.
- Promote it and raise the term in the consistent store.
- Fence: cut the old leader off, or make storage refuse its term.
- Point clients at the new leader: DNS, a proxy, or the router's leader key.
I promote the standby with the most WAL, raise the term, and fence the old leader so its late writes are refused.
| conflict rule | status | result |
|---|---|---|
| One home region per record | Approved | No conflicts: each record has one writer. |
| Merge with a CRDT: counters, sets, text | Mergeable data | Every edit survives, merged the same way on every copy. |
| Keep both versions, ask the app or user | Rare conflicts | Correct, but the app must handle siblings. |
| Last write wins by timestamp | Loss accepted | Drops one concurrent edit. Clock skew picks the winner. |
- Multi-leader: each region has a leader and accepts writes locally.
- Leaderless: any replica takes a write; quorums and read repair keep copies close.
- Both need a rule for two writes to the same key at the same time.
I avoid multi-leader writes unless users in several regions must write locally. Then I give each record a home region or merge conflicts with a CRDT.
| event | result | why it is safe | saved by |
|---|---|---|---|
| The leader crashes; standbys are asynchronous | The last commits may be missing. | Promote the standby with the most WAL. With a synchronous standby, nothing acknowledged is lost. | Sync standby |
| A network split leaves the old leader running | Two nodes think they lead. | The leader key expires first; the old leader's term is refused. | etcd lease |
| The only synchronous standby dies | Commits that wait for it block. | ANY 1 of 2 standbys: the other one acknowledges. | Quorum commit |
| A standby lags 30 seconds | Its reads are old. | The router reads replay_lag and skips that standby. | Lag check |
| A user reads right after their write | A replica may not have it. | The token makes the replica wait, or the read goes to the leader. | Token |
| A dead standby keeps its slot | WAL piles up on the leader. | max_slot_wal_keep_size caps it; an alert fires on slot lag. | Slot limit |
| A quorum replica is down during a write | A stand-in keeps a hint. | The hint is handed over later. Reads use R + W > N on home replicas. | Hinted handoff |
| step | add | it handles | move up when you see |
|---|---|---|---|
| 1 | One Postgres, with WAL archived for backups. | About 107,000 primary-key reads a second in the lab. Restore from backup after a loss. | A restore takes too long for your recovery time. |
| 2 | One asynchronous standby and a failover manager. | Failover within the leader key time to live. Loses only the commits inside the lag: 0.13 ms in the lab. | Reads load the primary. |
| 3 | Read replicas behind a router with lag checks and tokens. | Reads grow with each replica. Writes do not. | No acknowledged commit may be lost. |
| 4 | A synchronous standby: ANY 1 of 2. | Zero loss while one synchronous standby survives. Each commit waits one round trip. | The whole region can fail, or users are far away. |
| 5 | A standby in another region, asynchronous. | Survives a region loss. A synchronous one there adds about 40 ms per commit at 4,000 km. | Writes outgrow one leader: shard, as on the sharding sheet. |
Demand example: 1 billion reads a day is 1,000,000,000 / 86,400 ≈ 11,574 a second. Replica capacity is the lab measurement times the copy count. Light in fibre covers about 200 km per ms, so 4,000 km is 20 ms each way.
I start with one Postgres and archived WAL. Then I add a standby for failover, replicas for reads, and a synchronous standby when no commit may be lost.
0 of 10 known
A replica lags 50 ms. The user saves a change and reloads 10 ms later. What do they see, and how do you fix it?
Why does a sticky replica fix monotonic reads but not read-your-writes?
N = 3, W = 2, R = 2, and one replica is down. Do reads and writes work? Are reads fresh?
Why can a sloppy quorum return an old value even with R + W > N?
One synchronous standby is configured and it dies. What happens to commits?
You fail over to an asynchronous standby. What is lost?
The old leader wakes after a long pause and accepts writes. How do you stop it?
remote_apply makes reads on the standby fresh. Why not use it for every commit?
Two regions accept writes, and two users change the same field at the same time. What happens?
A standby dies, and its replication slot stays. What breaks next?
- lag, local
- A commit was visible on a standby on the same machine after 0.13 ms (median).
- sync commit
- 0.14 ms local, 0.36 ms with a synchronous standby on localhost: one more round trip.
- far standby
- 4,000 km is 20 ms each way in fibre: about 40 ms more per synchronous commit.
- off
- Can lose up to 600 ms of commits: 3 × the default wal_writer_delay.
- quorum
- N = 3, W = 2, R = 2 survives 1 replica down, with fresh reads.
- R = W = 1
- Of N = 3: stale in 6 of 9 cases when the read beats replication.
- reads
- About 107,000 primary-key reads a second on one Postgres.
Postgres 16 on an 8-core laptop; reads with 32 clients; commit latency over 500 commits. Simulations are seeded and repeat exactly.