Replication, Quorums & Read-Your-Writes
The Network Partition That Elected Two Leaders
Test your architecture intuition: Pitch a 7-axis solution, survive two aggressive reviewer objections, and inspect the staff-level Teacher Gold Answer.
Part 0. Start here
The problem: 500 lost orders and a stale checkout page
An orders database commits 5,000 transactions a second on its primary. It has two asynchronous replicas (copies that receive the primary's log after the commit, without the commit waiting for them), one in each of the other two Availability Zones: self-run PostgreSQL or MySQL, or RDS read replicas. The replicas serve reads and are the failover targets. During a heavy write burst, the primary's write-ahead log (WAL) reaches the replicas 100 ms after each commit (our example's number). At every instant of that burst, two things are true.
1. A failover loses acknowledged writes. Writes committed in the last 100 ms exist only on the primary:
- 5,000 commits a second × 0.1 s = 500 acknowledged transactions that no replica has yet.
If the primary dies now and a replica is promoted, 500 customers were told "OK" and their orders are gone. (Aurora is different: its replicas share one storage volume that already holds every commit, so an Aurora promotion loses nothing committed. Its lag only makes reads stale, which is the second problem.)
2. Users can't see their own writes. Anyone who saves and reloads before their write is applied on the replica their read lands on sees the old value. A checkout page that re-reads a just-saved shipping address from a replica ships to the old address. This happens on Aurora too: AWS says Aurora Replica lag is "usually much less than 100 milliseconds after the primary instance has written an update", and that it can rise "during periods where a large amount of write operations occur". A read inside that window, on the reader endpoint, is stale.
We want "OK" to mean the write survives the primary dying, and we want a user to see their own write on the next page. What do we change for each, and what does each change cost?
The big picture
Synthesizing vector architecture diagram...
What to notice: the write is committed when 2 of the 3 copies have it (the leader plus the fastest follower), and the OK carries a position. A read carrying that position may go to any copy that has applied it; if none has, it goes to the leader.
What you'll be able to do after this page
- Say what a replica keeps and name the three positions that matter: stored, committed and applied (Part 1).
- Say exactly when each replication mode acknowledges a write, what it loses in a failover, and what it does when followers are down (Part 2).
- Measure replication lag in its two parts, and say which part decides stale reads and which decides lost writes (Part 3).
- Give users read-your-writes and monotonic reads without sending every read to the leader (Part 4).
- Use N, W and R, and say what
W + R > Nguarantees and what it doesn't (Part 5). - Walk through a leader election after a crash and explain why no committed write is lost (Part 6).
- Explain what each side of a network partition can do, how a deposed leader is fenced, and how a database fails over (Part 7).
- Bring a lagging follower back from the log or from a snapshot, without starving everyone else (Part 8).
- Repair drifting copies, make deletes stick, and explain why replication is not a backup (Part 9).
- Choose single-leader, multi-leader or leaderless, and synchronous or asynchronous, on equal terms (Part 10).
- Map all of it to AWS services, and name the look-alikes that are not replication (Part 11).
You may have arrived from a step that relies on this: step 2.4 of the digital wallet loop (read-your-own-writes with min_version), step 1.3 of the message queue loop (Raft per log, commit on 2 of 3), step 1.6 of the key-value store loop ("when is a write done?", N = 3, W = 2, R = 2), or step 3.2 of the URL shortener loop (why an eventually consistent global table can't decide who owns an alias). This page is the "why" behind all four.
Part 1. What a replica keeps
Three copies, one truth: which copy is right when they disagree? Before any protocol, we need a picture of what each copy holds, so that "disagree" has a precise meaning.
One log, three copies
A replicated store keeps one ordered log of changes and copies it. The leader decides the order: each write becomes the next entry in its log. Each entry has a position (its index: PostgreSQL calls it an LSN, MySQL identifies transactions by GTID, Raft uses an index, Kafka an offset) and the term of the leader that created it (a counter that goes up with each new leader; Part 6). We draw an entry as term.index: 1.2 is the second entry, created by the term-1 leader.
"Replicated" means every copy has the same log in the same order. The copies differ only in how far along that log they are. So the questions on this page are all about positions: how far along is each copy, and which position does a write or a read wait for?
Stored, committed, applied
Every copy tracks three positions. They are different, and mixing them up is the source of most replication bugs.
| Position | Meaning | Who knows it | What it decides |
|---|---|---|---|
| Stored (received and flushed) | The entry is on this copy's disk | Each copy, about itself; the leader keeps each follower's latest acknowledged position (its match position) | Whether this copy would still have the entry if it became leader |
| Committed | The entry is safe: enough copies have it that no future leader can lack it | The leader decides it; followers learn it from the next message the leader sends (a heartbeat or an append) | When the client may be told OK |
| Applied | The entry has been executed into the copy's tables, so reads there see it | Each copy, about itself | What a read on this copy returns |
On any copy, applied ≤ committed as far as it knows, and it never applies an entry it hasn't stored. The three positions can be far apart: a copy can have an entry on disk for hundreds of milliseconds before it applies it (Part 3).
Synthesizing vector architecture diagram...
Snapshot S1, at t = 106 ms, from the reference implementation. All three copies have entry 1.2 on disk; only A has applied it. B will apply it when the next heartbeat tells it 1.2 is committed (110.2 ms); C won't until a long report query stops blocking its apply (400.0 ms). Same log, three different answers to a read.
Three families
| Family | Who accepts writes | How copies agree |
|---|---|---|
| Single-leader | One leader per log (or per shard) | Followers copy the leader's log in order. PostgreSQL, MySQL, Raft (etcd), Kafka, DynamoDB per partition |
| Multi-leader | Several leaders, each taking writes (often one per Region) | Leaders exchange changes and resolve conflicts. DynamoDB global tables (default mode), MySQL multi-source setups |
| Leaderless | Any replica, through a coordinator | The coordinator sends each write to all copies and waits for some; reads ask several and keep the newest version. Dynamo-style stores: Cassandra, ScyllaDB, Riak |
The page follows single-leader first, because most of the loops use it; Part 5 replays the same three copies without a leader, and Part 10 compares all three.
One exception needs no ordering at all: immutable data. A chunk of an object that is written once and never changed (the S3-like object storage loop, step 1.2) only needs enough copies; there is no "which version wins", because there is only one version. The mutable metadata record that points to the chunks is what needs a log.
What the WAL page already said
The Write-Ahead Log, fsync & Group Commit loop primitive's Part 10 has what "acknowledged" means on each engine (Kafka's acks=all, PostgreSQL's synchronous_commit, MySQL semi-sync, Raft, DynamoDB, Aurora) and what one replica's "on disk" costs. This page starts where that one stops: which copies wait, who leads, and which copy a read may use.
The example we follow
One key, addr:alice, on three replicas: A (AZ a), B (AZ b) and C (AZ c). Alice saves addresses from a laptop, a phone and a tablet; a checkout page, a report query, a repair job and a backup read them. We follow the key from its first write to its erasure. Every trace on this page comes from running a private reference implementation, a discrete-event simulator of these three replicas, not from working by hand.
| Setting | Tiny example | At real scale |
|---|---|---|
| Replicas | A (AZ a), B (AZ b), C (AZ c); A leads in term 1 | 3 or 5 replicas, one per AZ |
| Log | Entries term.index; the index is the position; addr:alice's version is the index of its last write | PostgreSQL LSN, MySQL GTID, Raft index, Kafka offset |
| Network | 0.2 ms one way between any two replicas (an assumption) | Measure your own cross-AZ round trip |
| Local sync | A and B 1.0 ms (the WAL page's value); C's volume is slower, 5.6 ms | Depends on the device (WAL page, Part 2) |
| Follower acknowledgement, from the leader's send | 0.2 + sync + 0.2: B 1.4 ms, C 6.0 ms (the WAL page's branch R values) | A cross-AZ round trip plus the follower's sync |
| Shipping | Raft-style: the leader sends entries at the same time as its own sync. The PostgreSQL rows in Part 2 ship only after the local flush | PostgreSQL's WAL sender reads only WAL already flushed to disk |
| Vote state | currentTerm and votedFor go to disk (the replica's own sync) before any vote request or reply is sent | Raft paper, Figure 2 |
| Heartbeat | Every 10 ms on the leader's tick; followers apply an entry once a heartbeat or append says it is committed | etcd: 100 ms |
| Election timeout | Randomized per server between 150 and 300 ms; in this run A drew 170, B 160 (later 290), C 240 | etcd: 1,000 ms, and at least 10 × the round trip |
| Leader read lease | 140 ms from the send time of the last heartbeat a majority acknowledged (Part 7) | Engine-specific |
| Apply on C | Blocked from 50 to 400 ms, and from 10,000 to 10,300 ms, by long report queries | PostgreSQL replay waits for conflicting standby queries up to max_standby_streaming_delay (default 30 s) |
| Client timeout | 2,000 ms | Application-specific |
The story has nine beats, each in its own Part: ways to say OK (Part 2), a stale checkout page (3), the fixes (4), the same replicas without a leader (5), A crashes mid-write (6), C is cut off (7), B comes back late (8), Alice is erased and a bad update reaches everyone (9), and choosing, with Alice in two Regions (10). Two side replays fork from the state after event 8, when all three copies hold "9 Elm St": the leaderless replay (rows Q1 to Q7, with their own clock, written q0, q3.0 and so on) and the two-Region replay (rows G1 to G3). The full event table is in Part 12.
What to remember from Part 1
- A replica is a copy of one ordered log, somewhere along it.
- "Stored", "committed" and "applied" are three different positions.
- Everything on this page is about which position a write or a read waits for.
Part 2. When to say OK: sync, async, semi-sync
At t = 0 the laptop saves Alice's first address, v1 "12 Oak St", which becomes entry 1.1 on A. At t = 1.2 ms, A loses power. Did Alice's write survive? It depends entirely on when A said OK, and on which copies had the entry by then. "We replicate" hides that choice; this Part makes it visible.
Trace: A dies at 1.2 ms
The same write under each mode, from the reference implementation. "OK at" is the time A would acknowledge if it lived; the last column is what happens if A dies at 1.2 ms and a follower is promoted.
| Mode | When it says OK | OK at | After A dies at 1.2 ms |
|---|---|---|---|
| Asynchronous | After the local sync; ships later | 1.0 ms | The OK already went out. Under the burst, WAL ships 50 ms after commit (a labelled assumption), so v1 would have left A at 51.0: it is only on A. B is promoted without it: an acknowledged write is lost |
| Raft-style quorum, W = 2 of 3 | After 2 copies (A plus the fastest follower) have it on disk; A sends in parallel with its own sync | max(1.0, 1.4) = 1.4 ms | No OK yet: A died first. But the copies were already in flight: B has v1 on disk at 1.2, C at 5.8. The new leader has v1, so the write survives although the client got no OK: its outcome is unknown to the client (Part 6) |
| All three, W = 3 of 3 | After every copy has it | max(1.0, 1.4, 6.0) = 6.0 ms | Same as the row above: kept, with no OK |
PostgreSQL synchronous, synchronous_commit = on with FIRST 1 (B) or ANY 1 (B, C) | After the standby has flushed it. The WAL sender ships only WAL already flushed locally, so the round trip starts after A's sync | 1.0 + 1.4 = 2.4 ms | The WAL left A at 1.0, after its flush: B has v1 on disk at 2.2, C at 6.8. Promoted B has v1; the client got no OK |
| MySQL semi-synchronous | After at least one replica acknowledges; a replica acknowledges "only after the events have been written to its relay log and flushed to disk". By default the source waits after syncing its binary log but before committing to the storage engine (AFTER_SYNC) | Not drawn: MySQL's documentation doesn't say whether events are sent before or after the binary log sync, so we don't time it | Depends on whether the events had left A |
Three lessons from one write:
- Asynchronous replication says OK before any other copy has the write. Any failover can lose the tail that hadn't shipped.
- A synchronous mode never loses a write it acknowledged to one failure, but it can keep a write it never acknowledged. The client that timed out must not assume "failed": a retry has to be safe to repeat, which is the Idempotency & Effectively-Once Processing loop primitive's subject.
- Latency depends on how the log is shipped. Raft-style logs send to followers while the leader syncs, so the commit waits for the slowest leg we need: max(local sync, the fastest follower's acknowledgement at W = 2). The slow follower C never enters the latency. PostgreSQL ships only flushed WAL, so its synchronous commit is local flush plus a round trip: 2.4 ms here, not 1.4. The WAL page's Part 10 has this latency rule with the same caveat.
When followers are down
The mode matters most when copies are missing. From the reference implementation, with the same timings:
| Situation | Raft W = 2 of 3 | PostgreSQL, standby list names only B (FIRST 1 (B)) | PostgreSQL ANY 1 (B, C) or FIRST 1 (B, C) | MySQL semi-sync |
|---|---|---|---|---|
| 1y: B is powered off | Commits when C acknowledges: 6.0 ms | Waits until B returns or someone changes the setting | 1.0 + C's 6.0 = 7.0 ms (FIRST moves to C, the next in priority) | Commits when C acknowledges: any one replica counts |
| 1z: B and C are both down | No commit; the client times out after 2,000 ms | Waits. PostgreSQL's documentation: such commits "may never be completed" if a required synchronous standby crashes | Waits | After rpl_semi_sync_source_timeout (10,000 ms by default) with no acknowledgement, falls back to asynchronous and commits locally |
The last cell is the dangerous one. Nothing fails, so nobody notices, but from that moment an OK means only "on the source's disk". The promise changed silently.
Semi-sync falls back to asynchronous after 10 seconds without a replica acknowledgement. What does the OK mean during those 10 seconds, and after them?
PostgreSQL's cancelled wait
A PostgreSQL commit that is waiting for its synchronous standby has already committed locally. If the client cancels the wait (or the connection is terminated), PostgreSQL can't un-commit it; it returns a warning instead:
textWARNING: canceling wait for synchronous replication due to user request DETAIL: The transaction has already committed locally, but might not have been replicated to the standby.
That commit is now visible on the primary and possibly on no standby. If the primary fails before the standby receives it, a failover loses it. This is the same lesson as the WAL page's "a write on one disk but not on a quorum is not kept", in a form that looks like success: treat a cancelled synchronous wait as an unknown outcome, never as a replicated commit.
How many writes an asynchronous failover can lose
For an asynchronous copy, the writes at risk are the ones committed on the leader but not yet stored on the copy we promote:
Formula 1:
Shipping lag is the distance between the leader's commit position and the position the replica has received and flushed, measured in time. It is not the replica's replay (apply) lag: a replica can have an entry safely on disk while it is still waiting to apply it (Part 3), and a failover to it keeps that entry. With the hook's numbers: 5,000 commits a second × 0.1 s = 500 writes.
What to remember from Part 2
- Async says OK before any other copy has the write; a failover can lose it.
- Quorum sync never loses an OK'd write to one failure, but a write that timed out may still be kept.
- Know what your mode does when followers are down: block, fall back to async, or carry on with the others.
Part 3. Replication lag and stale reads
At 150 ms Alice's checkout page shows "12 Oak St", the address she replaced 48.6 ms earlier. Nothing crashed; every copy is healthy. This Part is about the gap between a write committing and a copy showing it, which is where most replication bugs in the loops come from.
Trace: the stale checkout
Back to the main line, with v1 on all three copies (applied by 10.2 ms). The reference implementation:
| # | t (ms) | Event | A | B | C |
|---|---|---|---|---|---|
| 2 | 50.0 | A long report query starts on C. C keeps receiving and storing entries, but its apply stops behind the query | Oak St | Oak St | Oak St |
| 3 | 100.0 | Laptop writes v2 "9 Elm St" (entry 1.2). A syncs by 101.0; B's acknowledgement arrives at 101.4: committed, OK at 101.4. C has it on disk at 105.8, but can't apply it | Elm St | Oak St | Oak St |
| 4 | 110.0 | A heartbeat carries commit index 2; B applies v2 at 110.2 | Elm St | Elm St | Oak St |
| 5 | 150.0 | The checkout page reads through the reader endpoint and lands on C: "12 Oak St", stale, 150.0 − 101.4 = 48.6 ms after the OK | stale | ||
| 6 | 180.0 | Alice reloads; the read lands on B: "9 Elm St" | read | ||
| 7 | 200.0 | Next click lands on C: "12 Oak St" again. For Alice, time went backwards | stale | ||
| 8 | 400.0 | The report ends; C applies v2 at 400.0 | Elm St | Elm St | Elm St |
Synthesizing vector architecture diagram...
What to notice: C has v2 on disk from 105.8 ms, only 4.4 ms after the commit, yet both of its reads return the old address, because its apply is blocked until 400 ms. B, reached between them, shows the new one.
C has entry 1.2 on disk at 106 ms. Why does the read at 150 ms still show Oak St, and would a failover to C at 150 ms lose v2?
Shipping vs applying
| Lag | Definition | Decides | Our example |
|---|---|---|---|
| Shipping (receive) lag | Leader's commit position minus the replica's received-and-flushed position (or the time between them) | How many acknowledged writes an async failover loses (Formula 1) | 4.4 ms for v2 on C |
| Apply (replay) lag | Leader's commit position minus the replica's applied position; in time, now minus the commit time of the oldest committed entry not yet applied | How stale a read on that replica can be | 48.6 ms at the 150 ms read; peak 298.6 ms |
The metrics you get mostly measure the second, or both together. PostgreSQL's pg_stat_replication view, on the primary, has three: write_lag (the standby has written the WAL), flush_lag (written and flushed) and replay_lag (flushed and applied). Failover loss follows the flush gap; stale reads follow replay_lag, which PostgreSQL's documentation says, for an asynchronous standby, "approximates the delay before recent transactions became visible to queries". On a standby, comparing pg_last_wal_receive_lsn() with pg_last_wal_replay_lsn() separates the two. Aurora's AuroraReplicaLag metric and RDS's ReplicaLag are about when a replica shows the writer's changes, so they are read-staleness numbers.
Where lag comes from
| Cause | What happens | What to do |
|---|---|---|
| Network | Every entry crosses an AZ link; a slow or saturated link delays shipping | Watch the link; keep snapshot and repair streams throttled (Parts 8, 9) |
| A slow disk on the replica | It stores entries late (C's 5.6 ms sync) | Same volume class as the leader |
| Single-threaded apply | A replica applies one log in order, often on one thread, so a write burst the leader absorbed in parallel queues up | Parallel apply where the engine supports it; size the replica like the leader |
| A big transaction | Nothing in it is visible until all of it is applied | Batch large jobs into smaller transactions |
| A long query blocking replay | Replay would remove rows the query still needs, so it waits: our event 2. PostgreSQL waits up to max_standby_streaming_delay (30 s by default), then cancels the query | Accept cancelled queries, or set hot_standby_feedback on (off by default), which stops the conflicts but keeps dead rows on the primary longer (bloat) |
| Catching up after a restart | The replica starts behind by its whole downtime | Part 8 |
Measuring it and detecting a stuck replica
Position metrics tell you how far behind a replica is if it is still hearing from the leader. A replica that stopped receiving can look fine: MySQL's Seconds_Behind_Source compares the event being applied with its timestamp on the source, so when the applier has caught up with a slow receiver it shows 0, and MySQL's documentation warns that on a slow network it "often shows a value of 0, even if the replication receiver thread is late compared to the source"; if the receiver thread isn't running, it shows NULL. The robust check is a heartbeat row: once a second the primary writes the current time into a one-row table, and the monitor reads that row on each replica. Lag = now − the time in the row. It works whether the replica is slow, blocked or cut off, because a cut-off replica's row simply stops moving.
And lag is a distribution, not a number. "Usually much less than 100 milliseconds" and "typically within a second" are descriptions of the common case. They are never bounds: under a write burst, a blocked replay or a network problem, lag has no ceiling unless something enforces one (RDS's flow control for Multi-AZ DB clusters, or Aurora PostgreSQL's rds.global_db_rpo for a global database, both of which do it by slowing the writer). Design reads that tolerate the tail, and act on measured lag, like the paging library loop does when it widens its settle cutoff by the measured lag and moves reads to the writer above 5 s (steps R2.3 and R2.8).
Reads that go backwards
Events 5 to 7 show a second problem beyond staleness. The reader endpoint spreads reads across B and C, which are at different positions, so Alice saw new, then old. That breaks monotonic reads (never see an older state after a newer one). It is caused by load balancing across copies that differ, not by any single copy being slow.
What to remember from Part 3
- Lag is the distance between what the leader committed and what a replica stored or applied: measure both.
- Failover loss follows shipping lag; stale reads follow apply lag.
- Two reads that land on different replicas can go backwards in time.
Part 4. Read-your-writes and monotonic reads
Replay event 5: the checkout page's read at 150 ms is about to land on C, which hasn't applied v2. Which copy may answer? Neither "any replica" nor "always the leader" is a good answer. This Part gives the read enough information to choose.
The two promises
These are session guarantees: promises about what one user (one session) sees, not about the whole database. They are separate from transaction isolation levels, which are about concurrent transactions on one copy (Primitive #21: Isolation levels).
| Guarantee | Plain meaning | Broken in our story by |
|---|---|---|
| Read-your-writes | After I save something, my reads show it (or something newer) | Event 5: Alice saved Elm St at 101.4 and read Oak St at 150 |
| Monotonic reads | Once I've seen a state, I never see an older one | Event 7: Alice saw Elm St at 180, then Oak St at 200 |
Two more, named for completeness: monotonic writes (my writes apply in the order I made them) and writes-follow-reads (a write I make after reading something is ordered after what I read). A single-leader log gives both for free, because every write goes through one ordered log.
The mechanism: carry the position
The write's OK carries its commit position p. The read carries the highest position this session must see, and the router only uses a copy that has applied at least that far.
textON WRITE COMMIT: reply OK with p = the entry's position session.p_write = p -- on the client, or in a server-side session record ON READ: need = max(session.p_write, session.p_read) pick = the replica the load balancer chose if applied(pick) >= need: answer from pick else if another replica has applied(r) >= need: answer from r else wait up to a budget (say 20 ms) for pick to reach need else: answer from the leader session.p_read = the position of the copy that answered DECISIONS (money, permissions, consent): read the leader, always
The fixes, on equal terms
The reference implementation replays events 5 to 7 under six policies. In every row, v2 is position 2; C has applied position 1 until 400 ms; B has applied 2 since 110.2 ms.
| Fix | Read at 150 ms | Read at 200 ms | Guarantee | Cost | Breaks when |
|---|---|---|---|---|---|
| F1 Position token | The OK carried 2; the load balancer picked C (applied 1), so the router sends it to B (applied 2): "9 Elm St". Variant: wait up to 20 ms on C: still 1 at 170, so it goes to A | Same: B | Read-your-writes | A token per session; replicas must report their applied position | Every replica lags: all reads fall back to the leader |
| F2 Session high-water mark | As F1 | Alice's 180 ms read from B recorded position 2; C (1) is skipped: B, "9 Elm St" | Monotonic reads (and, with F1, read-your-writes) | The same token, updated on every read | The session record is lost |
| F3 Sticky routing to B | B: "9 Elm St" (applied at 110.2) | B: "9 Elm St" | Monotonic reads while B lives | Uneven load; a hot user pins a replica | Counter-row: stick Alice to C instead and she gets "12 Oak St" at both 150 and 200: monotonic, not read-your-writes. And when B fails, Alice moves to a copy that may be behind |
| F4 Leader reads for 1 s after a write | A: "9 Elm St" | A: "9 Elm St" | Read-your-writes while lag < 1 s | Leader load right after writes | Our 298.6 ms apply lag passes; a 30 s replay pause does not |
| F5 Server-side last-write position | Alice's phone reads at 160 with no local token; the server's session record for Alice holds 2, so it routes as F1: B | As F1 | Read-your-writes across devices | One more lookup per read | The session store is down: fall back to the leader |
| F6 Always read the leader | A: "9 Elm St" | A: "9 Elm St" | Read-your-writes; linearizable only with Part 7's read index or lease | The leader's CPU and I/O, and a cross-AZ round trip for readers in other AZs | The leader is the bottleneck, or it has been deposed without knowing (Part 7) |
Alice is pinned to replica B for her whole session. Is read-your-writes guaranteed?
Why not just read everything from the leader?
Engines and services
| Engine or service | The write's position | Checking it on a replica | Notes |
|---|---|---|---|
| PostgreSQL | pg_current_wal_lsn() on the primary after the commit (a WAL location at or past the commit record) | pg_last_wal_replay_lsn() on the standby, compared with the token; poll until it passes or the budget runs out | synchronous_commit = remote_apply makes a commit wait until the synchronous standbys have applied it: read-your-writes on those standbys, for the price of apply lag on every commit |
| MySQL | The transaction's GTID (the server can return it to the session) | WAIT_FOR_EXECUTED_GTID_SET(gtid_set, timeout) on the replica: waits until those transactions are applied; returns 0, or 1 on timeout | The timeout is the bounded wait |
| DynamoDB | None needed within one Region | ConsistentRead = true on GetItem, Query or Scan of a table or an LSI returns "the most up-to-date data" | GSIs and streams are always eventually consistent; eventually consistent reads cost half as much. With global tables in the default mode, a strongly consistent read can be stale if the item was last written in another Region |
| Caches in front (application caches, DAX) | The cache doesn't know positions | It serves what it holds | A cache in front of replicas hides the token check. In global tables, writes replicated from other Regions bypass DAX, which stays stale until its TTL expires |
The wallet loop (steps R2.3 and 2.4) uses exactly F1 with a min_version; the payments loop (step R2.5) and the nearby friends loop (step R1.7) use the rule "anything that decides reads the writer".
The read half, end to end through the layers
The write half of this trace, from the client down to "W of N have it", is in the WAL page's Part 10. Here is the other half: where staleness can enter a read, following the event-5 read at 150 ms.
Synthesizing vector architecture diagram...
What to notice: the position check is the only step that knows about Alice's write. A cache hit above it skips the check entirely, and the replica's apply queue below it (C's blocked apply) is exactly what the check protects against.
| Hop | Where staleness enters | What the token check sees at 150 ms |
|---|---|---|
| Client to app | A second device has no token of its own (F5 fixes it with a server-side record) | Laptop session: p = 2 |
| App cache or DAX | Serves an old value for its TTL; replicated writes bypass DAX | Nothing: the check never runs on a hit |
| Reader endpoint | Picks any replica, fresh or not | Picked C |
| Replica's applied position | The replica has stored but not applied (event 5) | C applied 1 < 2: without F1 the answer is "12 Oak St"; with F1 it goes to B |
| Replica's apply queue | A long query, a big transaction or single-threaded apply | C's queue is blocked until 400 ms |
| Response | Latency is the sum of the dependent hops, plus any bounded wait | With F1: one extra routing decision; with the 20 ms wait variant, 20 ms more before falling back to A |
What to remember from Part 4
- The write returns where it is; the read goes only to a copy that has got there.
- Sticky routing gives monotonic reads, not read-your-writes.
- Anything that decides money, permissions or consent reads the leader, with a read index or lease if it must be linearizable.
Part 5. Quorums: N, W and R
Take the same three copies and remove the leader. A read with R = 1 lands on C and returns "9 Elm St", although a newer address was acknowledged a moment ago. Without a leader, nobody orders the writes, so "enough copies" has to be counted on both sides: the writes and the reads. This Part is that counting rule and its limits.
The leaderless replay
We fork the story after event 8, when A, B and C all hold v2 "9 Elm St" (version 2). The copies are now Dynamo-style: a coordinator (a node beside A, in AZ a) sends every write to all N home replicas and waits for W acknowledgements; a read asks R replicas and keeps the answer with the highest version. The replay has its own clock (q0, q3.0 and so on) and its own versions, q3, q4 and q5. Replica answer times, from the coordinator's request to the reply: A 1.0 ms, B 1.4 ms, C 6.0 ms, stand-in D 1.2 ms.
From the reference implementation:
| # | q (ms) | Event | Result |
|---|---|---|---|
| Q1 | q0.0 | Laptop writes q3 "5 Ash Ct", N = 3, W = 2. The message to C is lost (a network blip) | A answers at 1.0, B at 1.4: OK at q1.4. C still has v2 |
| Q2 | q3.0 | Two reads. R = 1, answered by C | "9 Elm St": stale, although the write was acknowledged |
| q3.0 | R = 2, asked of A and C | Answers v3 from A and v2 from C, complete at q9.0 (the second answer is C's, at 6.0 ms): returns "5 Ash Ct". Read repair: the coordinator writes q3 back to C |
Synthesizing vector architecture diagram...
What to notice: the write was acknowledged by A and B, so any two replicas include at least one that has it. The R = 2 read met A and returned q3; an R = 1 read that met only C would have returned the old value.
The overlap rule
Formula 2:
If the number of copies a write waits for plus the number a read asks is more than the number of copies, every read set shares at least one copy with every acknowledged write's set, so the read meets that write. With N = 3, W = 2 and R = 2: 2 + 2 = 4 > 3. With R = 1: 2 + 1 = 3, not more than 3, which is why Q2's R = 1 read could miss it.
A single-leader log is the special case where W is a majority and the leader's own copy decides reads: every majority overlaps every other majority, which is why a new leader elected by a majority must meet every committed entry (Part 6).
| N, W, R | Overlap? | Writes survive | Reads survive | Write latency (A 1.0, B 1.4, C 6.0) | Read latency |
|---|---|---|---|---|---|
| 3, 1, 1 | No: 2 ≤ 3 | 2 copies down | 2 down | 1.0 ms | 1.0 ms |
| 3, 2, 1 | No: 3 ≤ 3 | 1 down | 2 down | 1.4 ms | 1.0 ms |
| 3, 2, 2 | Yes: 4 > 3 | 1 down | 1 down | 1.4 ms | 1.4 ms |
| 3, 3, 1 | Yes: 4 > 3 | none down | 2 down | 6.0 ms | 1.0 ms |
| 3, 1, 3 | Yes: 4 > 3 | 2 down | none down | 1.0 ms | 6.0 ms |
Latency is the W-th fastest replica for a write and the R-th fastest of those asked for a read (in the table, the reads ask the fastest replicas first). When the coordinator is itself a replica, a write waits for max(its own sync, the (W − 1)-th fastest other replica), the same rule as the Raft leader in Part 2. With 5 replicas and W = 3, a write waits for the second-fastest of the four others: more copies can make a write faster at the tail, because the slowest one is ignored, while costing more bytes and storage (Part 7 has the 3, 4 and 5 table).
Most stores let you choose per request: Cassandra's ONE, QUORUM and LOCAL_QUORUM (a quorum of the replicas in the local data center), as the key-value store loop does in step 2.1.
Where the overlap stops
W + R > N guarantees one thing: in normal operation, a read meets the latest acknowledged write. It does not make the store linearizable (every read sees every write that finished before it began, as if there were one copy). Three cases break the simple promise:
| Case | What happens |
|---|---|
| Concurrent writes | Two writes to the same key at the same time each reach W copies. Versions or timestamps decide which wins, and under last-writer-wins one of them silently disappears. The key-value loop (step 2.2) keeps both as siblings with vector clocks instead |
| A write that failed partway | A write that reached 1 copy and then reported an error is neither acknowledged nor rolled back. A later read may or may not see it, depending on which copies it asks, and read repair can then spread it |
| A quorum read, then a write | Reading a value with R and writing a new one with W is not compare-and-set: two clients can read the same value and both write. An atomic condition needs a consensus round per operation; Cassandra's lightweight transactions (conditional writes with IF) use Paxos for this (key-value loop, step 2.3) |
The consistency ladder
"Is a quorum store linearizable?" is the question behind all three rows. One table, strongest first, each with the mechanism on this page that provides it:
| Level | Plain meaning | Provided here by |
|---|---|---|
| Linearizable | Behaves like one copy: a read sees every write that finished before it started | A leader that confirms it is still the leader: read index or lease (Part 7); consensus per operation |
| Sequential | Everyone sees the same order of writes, but maybe late | (named only) |
| Causal | If one write could have influenced another, everyone sees them in that order | (named only) |
| Session guarantees | Read-your-writes and monotonic reads, for one session | Position tokens (Part 4) |
| Eventual | Copies converge if writes stop | Any replica read; read repair and anti-entropy (Part 9) |
W + R > N on its own sits below linearizable: it gives "a read overlaps the latest acknowledged write" when there are no concurrent or failed writes, and nothing when the quorum is sloppy (next).
Sloppy quorums and hinted handoff, briefly
What if home replicas are down when a write arrives? A strict quorum counts only the N home replicas, so the write fails. A sloppy quorum lets a stand-in node accept the write with a hint naming the home replica it belongs to, and counts the stand-in toward W; the stand-in hands the write over when the home replica returns (hinted handoff). The 2007 Dynamo paper and Riak do this. Cassandra stores hints too, on the coordinator, but a hint doesn't count toward the consistency level (except at the special level ANY), so Cassandra's quorums stay strict. The key-value loop warns about exactly this in step 2.4 and its mistake list (R2.9 #5): a sloppy quorum mistaken for strong consistency.
N = 3, W = 2, R = 2, and a read at q40 still returns the old address. Which rule was broken?
Snapshot S7, after Q5 (q40): A holds q4 "7 Birch Ave"; B and C hold q3 "5 Ash Ct"; the stand-in D holds hints of q4 for B and C.
What to remember from Part 5
- Overlap needs W + R > N counted over the same N home replicas.
- Overlap means a read meets the latest acknowledged write in normal operation, not linearizability.
- A sloppy quorum trades the overlap for availability until the hints are delivered.
Part 6. The leader crashes: elections
Back on the main line. At 1,003.0 ms A loses power in the middle of a write: its entry 1.3 reached C but not B. Someone must take over, and the new leader must hold every committed write, including ones it can't tell were committed. This Part is how Raft does that; the same ideas (a term, one vote per term, "only the most up-to-date may win") appear in every consensus-backed failover.
Trace: A dies with index 3 half-sent
| # | t (ms) | Event |
|---|---|---|
| 9 | 990.0 | A's last heartbeat to reach both followers (received at 990.2) |
| 10 | 1,000.0 | The laptop writes v3 "3 Pine Rd", entry 1.3. The append to B is lost. C receives it at 1,000.2 (which also resets C's election timer) and has it on disk at 1,005.8; its acknowledgement would arrive at 1,006.0 |
| 11 | 1,003.0 | A loses power. A had 1.3 on its own disk since 1,001.0, but only 1 of 3 copies had it: not committed, no OK. The laptop will time out at 3,000.0 |
Snapshot S2 (t = 1,003.0, from the reference implementation):
| A | B | C | |
|---|---|---|---|
| Log | 1.1 1.2 1.3 | 1.1 1.2 | 1.1 1.2 1.3 (1.3 on disk at 1,005.8) |
| Knows committed / applied | 2 / 2 | 2 / 2 | 2 / 2 |
| Term, vote | 1, A | 1, A | 1, A |
| Reads | dead | "9 Elm St" | "9 Elm St" |
Terms, votes, and why votes are written down
Raft's election rules, as far as our story needs them:
- Terms are a logical clock: each election starts a new, higher term, and each term has at most one leader. Every message carries its sender's term. A node that sees a higher term adopts it and becomes a follower; a message from a lower term is refused.
- Randomized timeouts: a follower that hears nothing from a leader for its election timeout (drawn at random, 150 to 300 ms in the Raft paper's example) becomes a candidate: it increments its term, votes for itself and asks the others for votes. Randomness makes it unlikely that two candidates split the votes; a split just means another round with new random timeouts.
- One vote per term, written to disk before it is sent. A node's current term and its vote are persistent state, updated on stable storage before it responds. Otherwise a node could vote, crash, restart with no memory of it, and vote again for someone else in the same term: two leaders in one term.
- The election restriction: a voter refuses a candidate whose log is less up-to-date than its own. The Raft paper's rule: "If the logs have last entries with different terms, then the log with the later term is more up-to-date. If the logs end with the same term, then whichever log is longer is more up-to-date."
textON ELECTION TIMEOUT (no leader heard for my random timeout): term = term + 1; votedFor = me sync term and votedFor to disk -- before anything is sent send RequestVote(term, my last entry's term and index) to all ON RequestVote(candidate, term, lastTerm, lastIndex): if term > myTerm: myTerm = term; votedFor = none; become follower upToDate = lastTerm > myLastTerm or (lastTerm == myLastTerm and lastIndex >= myLastIndex) grant = (term == myTerm) and votedFor in (none, candidate) and upToDate if grant: votedFor = candidate; reset my election timer sync myTerm and votedFor to disk -- before replying reply(myTerm, grant) A CANDIDATE WITH VOTES FROM A MAJORITY (itself included) IS LEADER FOR THAT TERM
Why the restriction keeps every committed entry: a committed entry is on a majority, a winner needs votes from a majority, and any two majorities share a node. That shared node has the entry and refuses any candidate whose log lacks it.
Trace: who wins
From the reference implementation (the arithmetic shown in brackets):
| # | t (ms) | Event |
|---|---|---|
| 12 | 1,150.2 | B's timer fires (last heard 990.2 + its 160 ms timeout): term 2, votes for itself, syncs term and vote (1.0 ms) |
| 1,151.2 | B sends RequestVote(term 2, last entry 1.2) to A (dead) and C | |
| 1,151.4 | C receives it. C's last entry is 1.3, same term as B's, and longer: C is more up-to-date, so it refuses. It adopts term 2 and persists it on its slow disk (5.6 ms) before replying | |
| 1,157.2 | The refusal reaches B (1,151.4 + 5.6 + 0.2). B has 1 vote of 3: no leader. A refusal doesn't reset C's timer | |
| 13 | 1,240.2 | C's timer fires (last heard 1,000.2 + 240): term 3, votes for itself, syncs (5.6 ms) |
| 1,245.8 | C sends RequestVote(term 3, last entry 1.3) | |
| 1,246.0 | B receives it: term 3 is higher, and C's log is at least as up-to-date as B's: B grants, syncs (1.0 ms) | |
| 1,247.2 | The grant reaches C: 2 of 3 votes. C leads term 3. It sets nextIndex (the next entry to send each follower) to its last index + 1 = 4 | |
| 14 | 1,247.2 | C appends a no-op entry 3.4 and starts syncing it (5.6 ms, done at 1,252.8). It sends B an append: "previous entry 1.3, then 3.4" |
| 1,247.4 | B refuses: it has no entry 3 | |
| 1,247.6 | The refusal reaches C, which moves nextIndex[B] back to 3 and resends 1.3 and 3.4 | |
| 1,248.8 | B has both on disk; its acknowledgement reaches C at 1,249.0 | |
| 1,252.8 | C's own sync of 3.4 finishes: 3.4 is on {B, C}, a majority, and is from C's own term: index 4 is committed, and index 3 with it | |
| 15 | 1,500.0 | A restarts with 1.1 to 1.3. C's next heartbeat (term 3, "previous entry 1.3", then 3.4) matches A's log; A adopts term 3, stores 3.4 (at 1,501.2) and follows |
B's timer fired first, yet C leads, because B's log was missing an entry C had. That is the restriction doing its job: 1.3 might have been committed as far as anyone alive knew.
Synthesizing vector architecture diagram...
What to notice: the first candidate loses because its log is shorter; every vote waits for a disk sync; and the commit waits for C's slow disk, which is now the leader's disk.
The no-op and the commitment rule
When C wins, it has entry 1.3 on its disk and knows a majority ({B, C}) will soon have it. Can it call 1.3 committed as soon as B stores it? Raft says no:
textLEADER, AFTER A FOLLOWER ACKNOWLEDGES: for N from my last index down to commitIndex + 1: if log[N].term == currentTerm and a majority (me included) has stored N: commitIndex = N -- every earlier entry is committed with it break ON WINNING AN ELECTION: append a no-op in the new term and replicate it
C holds entry 1.3, and so does B after 1,248.8. Why doesn't C call index 3 committed as soon as a majority has it, instead of waiting for its no-op?
Snapshot S3 (t = 1,252.8):
| A | B | C | |
|---|---|---|---|
| Log | 1.1 1.2 1.3 (dead) | 1.1 1.2 1.3 3.4 | 1.1 1.2 1.3 3.4 |
| Committed | (dead) | 4, learned at the next heartbeat, 1,260.2 | 4 |
| Applied, value | (dead) | 3 at 1,260.2: "3 Pine Rd" | 3 at 1,252.8: "3 Pine Rd" |
| Term, vote | 1, A | 3, C | 3, C (leader) |
The write that committed anyway
Alice's laptop never got an OK for v3: A died before it could send one, and the laptop times out at 3,000.0. Yet v3 is committed at 1,252.8 and is Alice's address everywhere. A client that treats a timeout as "failed" and writes something else, or retries a non-idempotent operation ("add $10"), gets it wrong. A timeout is an unknown outcome: retries must be idempotent, which the Idempotency & Effectively-Once Processing loop primitive builds.
Two costs of the election show in the trace:
- The write gap: from the crash to the first commit, 1,252.8 − 1,003.0 = 249.8 ms in which nothing could commit, including the vote syncs and a failed first round. At etcd's defaults (a 1,000 ms election timeout), expect seconds.
- The new leader's disk now sets the latency. Every commit in term 3 waits for C's own 5.6 ms sync: max(5.6, B's 1.4) = 5.6 ms. An election can move leadership onto the slowest disk; some systems prefer or transfer leadership to a better node afterwards.
What to remember from Part 6
- Only a candidate whose log is at least as up-to-date as a majority's can win, so every committed entry survives.
- A new leader commits a no-op in its own term; older entries become committed through it.
- A write that timed out can still commit: retries must be idempotent.
Part 7. A partition and a split-brain attempt
The WAL page's branch R (ii) showed the short version: a leader cut off from its followers syncs writes to its own disk, the others elect a new leader, and when the old leader rejoins it deletes those writes. Here we slow that down. At 2,001 ms C's whole AZ is cut off from A and B, while Alice's tablet, in AZ c, can still reach C. C still believes it leads.
Trace: C keeps writing
| # | t (ms) | Event |
|---|---|---|
| 16 | 2,000.0 | C sends a heartbeat; A and B receive it at 2,000.2 and acknowledge. A majority acknowledged a heartbeat sent at 2,000.0, so C's read lease runs to 2,000.0 + 140 = 2,140.0. At 2,001.0 C's AZ is cut off from A and B |
| 17 | 2,010.0 | The tablet writes v5 "1 Cedar Ln" to C: appended as 3.5, on C's disk at 2,015.6. No follower answers, so it is never committed; the tablet times out at 4,010.0 |
| 18 | 2,140.0 | C's read lease ends. The earliest another node could time out is 2,000.2 + 150 (the shortest possible timeout) = 2,150.2: a 10.2 ms margin |
| 19 | 2,170.2 | A's timer fires (2,000.2 + 170): term 4, synced by 2,171.2; B grants after its own sync; the grant arrives at 2,172.6: A leads term 4. It sets nextIndex[C] = 5 and appends no-op 4.5, committed on {A, B} at 2,174.0 |
| 20 | 2,200.0 | The phone writes v6 "22 Maple Dr" (4.6) to A: committed and OK at 2,201.4 |
| 21 | 2,250.0 | Reads at C, three ways (below) |
| 22 | 3,000.0 | The link is restored (next section) |
Two leaders now exist, in different terms. Snapshot S4 (t = 2,201.4):
Synthesizing vector architecture diagram...
What to notice: both panels have a leader, but only "Majority: A and B" can commit. C's entry at index 5 has term 3; A's entry at the same index has term 4. That difference is what settles it when the link returns.
What each side can do
| Minority (C) | Majority (A and B) | |
|---|---|---|
| Elect a leader | No: 1 vote of 3 | Yes, in a higher term (term 4) |
| Append | Yes, to its own log | Yes |
| Commit | Never: it needs 2 of 3 | Yes |
| Answer reads | Only stale ones, if it answers without checking | Yes, from the new leader |
This is the drill The Network Partition That Elected Two Leaders in log terms. Page 03's Part 8 gave the lock-service consequence (a minority can't grant, release or expire a lock). Inside the log: in a 5-node cluster split 2 and 3, the 2-node side can't commit any client write, because a commit needs 3 of 5 copies. Its old leader keeps appending entries that time out; the 3-node side elects a leader in a higher term and carries on; when the partition heals, the old leader sees the higher term, steps down, and its uncommitted entries are deleted. Nothing is merged: the majority's log wins.
Terms fence the old leader
Every message carries its sender's term. When the link returns, the first message C sees from A has term 4, higher than C's 3: C becomes a follower at once. Any message C sends with term 3 is refused by A and B with "my term is 4". This is term fencing, and it covers everything that goes through the log. It doesn't cover anything the old leader did outside the log (an email it sent, a file it wrote to S3): those need a fencing token checked by that other system, which is the Leases, Fencing Tokens & Distributed Locks loop primitive's subject.
Reads on a leader that may be deposed
At 2,250 ms a client asks C for Alice's address. From the reference implementation:
| How C answers | Result at 2,250 ms |
|---|---|
| No check: read local applied state | "3 Pine Rd": stale. Alice's real address, "22 Maple Dr", was committed by the term-4 leader at 2,201.4 |
| Read index: record the commit index (4), confirm leadership with a heartbeat round to a majority, wait until applied, answer | The heartbeat round never gets a majority. The read waits, then fails at 3,000.2 when C steps down on hearing term 4. The client retries at the new leader |
| Lease: answer locally only while the lease lasts | Refused: the lease ended at 2,140.0 |
textREAD INDEX (a round trip per read, or per batch of reads): r = commitIndex send a heartbeat to every follower; wait for acks from a majority (me included) if I stepped down meanwhile: fail, redirect to the leader wait until applied >= r; answer from local state LEASE READ (no round trip, but relies on clocks): lease_start = send time of the last heartbeat a majority acknowledged answer locally only while now < lease_start + lease choose lease < minimum election timeout - a margin for clock-rate difference
Why time the lease from the heartbeat's send: a follower resets its election timer when it receives the heartbeat, which is after the send, so the earliest any follower can start an election is (receive time + minimum timeout), which is later than (send time + minimum timeout). Ending the lease 10.2 ms before that leaves room for clocks that run at slightly different speeds. The Raft paper notes a lease "would rely on timing for safety"; how much margin clock behaviour needs is page 03's Part 3.
| Read index | Lease | |
|---|---|---|
| Cost | A heartbeat round per read, or per batch of reads | None per read |
| Depends on | Nothing but messages | Clocks running at close to the same rate; bounded pauses |
| On a deposed leader | Can't get a majority: fails | The lease has expired before anyone else can lead: refuses |
When the link returns
At 3,000.0 the link is restored. A's next heartbeat to C, in term 4, carries "previous entry 3.4, then 4.5 and 4.6".
Synthesizing vector architecture diagram...
What to notice: the previous entry matches, so C doesn't need to back up; the conflict is at index 5, where the terms differ, and C deletes from there. The tablet's write, never committed, is gone, and the tablet had never been told OK.
The truncation rule: when a follower accepts an append, any of its entries that conflict with the leader's (same index, different term) are deleted, together with everything after them, and the leader's entries are appended. Only uncommitted entries can ever be deleted this way, because a leader always holds every committed entry (Part 6).
Snapshot S5 (t = 3,010): all three copies hold 1.1 1.2 1.3 3.4 4.5 4.6, term 4, commit 6, and read "22 Maple Dr". C had 4.5 and 4.6 on disk at 3,005.8 and applied 4.6 then.
The tablet's "1 Cedar Ln" sat on C's disk for almost a second. When the link returns, is it lost, kept, or merged with the phone's "22 Maple Dr"?
3 vs 4 vs 5 replicas
| Replicas | Majority | Failures tolerated | A commit waits for | Storage and bytes |
|---|---|---|---|---|
| 3 | 2 | 1 | The faster of 2 followers | 3 copies |
| 4 | 3 | 1 (no better than 3) | The second-fastest of 3 followers | 4/3 of three copies |
| 5 | 3 | 2 | The second-fastest of 4 followers | 5/3 of three copies |
Even counts buy nothing: 4 tolerates one failure, like 3, at more cost. Five tolerates two, which matters most during maintenance, when one node is already down on purpose; the message queue loop (step R1.8) weighs exactly this. That is drill 11's second question: 3 nodes are cheaper and a commit waits for one follower, but one planned restart plus one failure stops writes.
When a database fails over
Database failover is the same election, run by a manager instead of by the replicas themselves:
- Detect that the primary is gone (health checks, a consensus-backed manager, or the managed service).
- Fence the old primary before promoting: stop it, revoke its access, or raise an epoch its writes must carry (page 03). Skipping this is how split brain happens.
- Choose the most caught-up replica: the one with the highest stored position.
- Promote it, and repoint the writer endpoint.
- Bring the old primary back as a replica only after throwing away its unreplicated tail. PostgreSQL's
pg_rewinddoes this: it synchronizes a data directory with another copy "after the clusters' timelines have diverged", which is the old-primary case. - Asynchronous replication loses the tail between the last shipped position and the crash (Formula 1).
| Service or setup | What is promoted | What a failover can lose | In-flight work |
|---|---|---|---|
| RDS Multi-AZ DB instance | The synchronous standby, which can't serve reads | Nothing acknowledged | Open transactions are aborted; commits with no OK received have an unknown outcome; connections and prepared statements drop; clients reconnect through the endpoint's DNS name (mind client DNS caching) |
| RDS Multi-AZ DB cluster | A reader: RDS picks the one "which has the most recent change record" (semi-synchronous: at least one reader acknowledged every commit) | Nothing acknowledged | Same |
| Aurora | An Aurora Replica: all share the storage volume that holds every commit | Nothing committed | Same; AWS describes "a brief interruption" |
| Self-run PostgreSQL with a failover manager | The replica with the highest position | Async: the unshipped tail; sync: nothing acknowledged, if the manager fences and picks correctly | Same, plus whatever your proxy or DNS does |
| ElastiCache (Valkey, Redis OSS) without durability | A replica, kept in sync asynchronously | Acknowledged writes in the unshipped tail | Clients reconnect to the new primary |
| Kafka (MSK) | A new partition leader from the in-sync replicas, unless unclean election is on | With unclean election: acknowledged messages | Producers retry; use idempotent producers |
Kafka's unclean leader election is the version of "promote a replica that lacks committed entries". Kafka's leader must normally come from the in-sync replicas (those caught up within replica.lag.time.max.ms). If every in-sync replica is lost, unclean.leader.election.enable=true lets an out-of-sync one lead, and the messages it never received are gone even though producers were told they were acknowledged. Apache Kafka's own default is false, but MSK's default configuration sets it to true for clusters without tiered storage (false with tiered storage). Set it to false for data you promised to keep, and accept that the partition is unavailable until an in-sync replica returns.
Failover between Regions (detection, routing, RPO and RTO decisions) is the Multi-Region Failover loop primitive's subject (coming).
What to remember from Part 7
- A minority can append but never commit; its conflicting entries are deleted when it rejoins.
- A deposed leader can serve stale reads: use a read index, or a lease that ends before any election can finish.
- Failover promotes the most caught-up copy and fences the old one first; an unclean election or an async promotion loses committed writes.
Part 8. Catching up: log or snapshot
At 9,000 ms B comes back from a host failure. It needs entry 7 next, but the leader deleted entries 1 to 10 hours ago (in toy time, 3 seconds ago) to keep its log short. The leader can't send what it no longer has. This Part is how a follower catches up, and what the catching up costs everyone else.
Trace: B returns to a compacted log
| # | t (ms) | Event |
|---|---|---|
| 23 | 4,000.0 | B's host fails; its disk survives. B's log ends at index 6 |
| 24 | 4,100 to 8,000 | Other customers' writes (not shown) commit as indexes 7 to 14 on {A, C}. With B gone, every commit waits for C's slow disk: max(A's 1.0, C's 6.0) = 6.0 ms per commit instead of 1.4 |
| 25 | 6,000.0 | A snapshots its state at index 10 and deletes entries 1 to 10: its log now starts at 11 |
| 26 | 9,000.0 | B restarts. A's heartbeat at 9,000.0 reaches B; B's reply at 9,000.4 shows its last entry is 4.6. nextIndex[B] is 7, before A's first entry (11), so A sends its snapshot at index 10 at 9,000.4 |
| 9,801.6 | B has installed the snapshot (a 400 MB snapshot at a throttled 500 MB/s takes 800 ms; both sizes are toy assumptions) | |
| 9,803.0 | A sends entries 11 to 14; B has them on disk. B is a full member again, and commits drop back to 1.4 ms |
Synthesizing vector architecture diagram...
What to notice: the log is always the first choice, because it sends only what is missing. A snapshot is the fallback when the log no longer reaches back far enough, and it is only useful if the log continues exactly where it ends.
Snapshot S6 (t = 9,803.0):
| A (leader, term 4) | B | C | |
|---|---|---|---|
| Log | snapshot@10, then 4.11 to 4.14 | snapshot@10, then 4.11 to 4.14 (entries 1 to 6 replaced) | 1.1 … 4.14 |
| Committed / applied | 14 / 14 | 14 / 14 (commit learned with entries 11 to 14 at 9,802.0; applied once stored, 9,803.0) | 14 / 14 |
addr:alice | "22 Maple Dr" | "22 Maple Dr" | "22 Maple Dr" |
Log or snapshot
- Catch-up by log: the leader finds the last entry the follower has that matches its own (the backtracking of event 14, one step back per refusal or straight to the follower's hint) and sends everything after it. Cheap: it sends only the missing entries.
- Catch-up by snapshot: when the entries the follower needs have been compacted away, the leader sends a copy of its state at some position s, and then the log from s + 1. Expensive: the whole state, however little was missing.
The validity rule: a copy seeded from a snapshot at position s may follow the stream only if the stream still has everything from s + 1 onward (its first entry is at or before s + 1). Physical replication engines enforce it: a Raft leader sends a snapshot automatically when it lacks the entries, and a PostgreSQL standby that asks for WAL the primary has already removed gets an error ("requested WAL segment … has already been removed") instead of silently skipping it. The silent version of this gap, a change-stream consumer or a hand-built tool that starts from a snapshot and a later stream position, is the Change Streams & the Transactional Outbox loop primitive's Part 7.
B needs entry 7; the leader's log starts at 11. What does the leader send, and what does sending it cost the other follower, C?
Throttle every bulk transfer that shares a leader's disk and network with live replication: PostgreSQL's pg_basebackup --max-rate ("useful to limit the impact of pg_basebackup on the server"), Cassandra's stream_throughput_outbound (Cassandra's own configuration notes that streaming "can lead to saturating the network connection and degrading rpc performance"), and whatever your Raft library offers for snapshot sends.
The other side of the choice: a leader could keep its log longer so that followers always catch up by log. But the log's retention is pinned by its slowest reader, follower, backup shipper or change-stream consumer alike, and a reader that never returns fills the disk (the WAL page's Part 8). Engines bound it (a snapshot threshold, a slot size cap) and fall back to snapshots.
New members join the same way, and they should not vote until they have caught up: a node that counts toward the majority but has an empty log makes every commit wait for it or for someone else. etcd can add new members as learners (non-voting) for this reason. How membership changes themselves stay safe is beyond this page.
What to remember from Part 8
- A follower catches up from the log while the log still has what it needs; otherwise from a snapshot.
- A snapshot is valid only if the stream continues from exactly where it ends.
- A snapshot competes with live replication for the leader's network and disk: throttle it.
Part 9. Repair, deletes and what replication can't protect
Back in the leaderless replay, at q1,400: Alice's account was erased an hour ago (in toy time, 1.2 seconds ago), and her address "7 Birch Ave" is back on every copy. Nobody wrote it. A repair process copied it from a replica that had missed the delete. This Part covers how copies are repaired, why deletes need special care, and the failures that replication copies instead of preventing.
Trace: Alice's erased address comes back (Q7)
From the reference implementation, continuing the leaderless replay (all three copies hold q4 after Q6):
| q (ms) | Event | Tombstone grace 1,000 q-ms | Grace 2,000 q-ms |
|---|---|---|---|
| q150 | B goes down | ||
| q200 | Alice is erased: the coordinator writes a tombstone q5 (a delete marker with its own version) with W = 2. A stores it at q201.0 (the coordinator sits beside A), C at q205.8: OK at q206.0. B is down, so a hint for B is kept on node D | same | same |
| q300 | D is replaced: its disk, and B's hint, are gone | same | same |
| q1,201 and q1,206 | A's and C's compactions purge the tombstone: it has been older than the grace period since q1,201.0 and q1,205.8 | Purged | Kept (eligible only from q2,201.0) |
| q1,400 | B returns, still holding q4 "7 Birch Ave". Anti-entropy compares Merkle trees, A with B and then C with B: 1 of 4 key ranges differs | A and C have no record of addr:alice; B has q4: repair copies q4 to A and C. Alice's address is back everywhere | A and C have tombstone q5, newer than B's q4: repair copies the tombstone to B. Alice stays erased |
Read repair and anti-entropy
Read repair is what Q2 showed: a read that asks R ≥ 2 copies compares their versions, returns the newest, and writes it back to any copy it asked that was behind. Stores often ask one copy for the full value and the others only for a digest (a hash), which saves cross-AZ bytes; only a mismatch triggers full reads. Read repair only fixes keys somebody reads, and in Cassandra it is best-effort.
Anti-entropy repairs everything else. Each replica builds a Merkle tree: hashes of key ranges, then hashes of those hashes, up to one root. Two replicas compare roots; if they match, the ranges are identical and nothing moves. If they differ, they compare children and descend only where hashes differ, then copy just those ranges.
Synthesizing vector architecture diagram...
What to notice: the roots differ, but "Replica A" and "Replica B" agree on the right half (8e5e), so it is skipped; on the left half only the range "a to c" differs, and it is the only one copied. The labels show the first characters of each hash from the reference implementation's run with the 1,000 q-ms grace.
Both kinds of repair use the same cross-AZ links and replica disks as foreground traffic: building trees costs reads, and a large difference means a large transfer. Schedule and throttle them like snapshot transfers (Part 8).
A delete is a versioned write
A missing value looks exactly like a write that never arrived. So a delete can't simply remove the key: it writes a tombstone with a version, and repair treats it like any other newer version. The tombstone must stay until every copy has it. In a leaderless store, that means longer than any replica can be away, plus one full repair cycle. Cassandra calls this grace period gc_grace_seconds, 864,000 s (10 days) by default, and its documentation ties it to hints and repair: hints older than the grace period are not replayed, and hints for a down node are collected only for max_hint_window (3 hours by default). The key-value store loop (step 2.6) has the same rule: run a full repair on every replica more often than the grace period, or deleted data comes back.
A single-leader log avoids this particular trap: the delete is an entry in the log, and a returning follower must replay the log or install a snapshot taken after it (Part 8).
Which copy may a repair read?
A repair or compare job must read the leader, or a copy whose applied position is at or past the delete. Event 27 on the main line shows the window: at 10,000 ms a report blocks C's apply until 10,300 ms, and DELETE addr:alice (entry 15) commits at 10,001.4 on {A, B}. A applies it at 10,001.4, B at 10,010.2, but C still returns "22 Maple Dr" until 10,300.0. A job that reads C during that window and "fixes" another system from it brings the address back. The Change Streams & the Transactional Outbox loop primitive (Part 8, its event 28) shows exactly that job, and the tombstone that saves its consumers.
Snapshot S8 (t = 10,100): A and B have applied the delete (no addr:alice); C has stored entry 15 but applied only 14, and still reads "22 Maple Dr".
Backups remember
Replication applies the delete everywhere within milliseconds. Backups do not:
| Where a deleted value still lives | For how long |
|---|---|
| Lagging replicas | Until they apply the delete (event 27: 300 ms) |
| Snapshots and base backups | Until they expire |
| Point-in-time restore windows (PITR) | The whole window: a restore to a time before the delete brings the value back (event 28: last night's snapshot and the PITR window still contain "22 Maple Dr") |
| Exports, analytics copies, caches, search indexes | Until they are rebuilt, expire or receive the delete |
So an erasure promise ("we delete your data within 30 days") must name every one of those windows, or encrypt each user's data with a per-user key and delete the key (crypto-shredding), which makes every remaining copy unreadable at once. The YouTube loop (step 3.4) and the Google Drive loop (step R3.11) make the same point.
One bad update everywhere
Event 29, from the reference implementation: at 12,000.0 a buggy migration sets every address to "" in one transaction (entry 16). It commits at 12,001.4; A applies it at 12,001.4, B and C at 12,010.2. Every copy agrees on the wrong value within about 10 ms. Replication did its job perfectly: it copied the mistake as fast as it copies any write.
The recovery needs a copy from before the mistake: a point-in-time restore to t = 11,999 into a new cluster, then copying the correct values back (in real life, hours later); or a delayed replica that applies the log a fixed time behind (MySQL's SOURCE_DELAY: a transaction is "not executed until at least N seconds later than its commit"; PostgreSQL's recovery_min_apply_delay), which can be stopped before it applies the bad entry.
Three replicas in three AZs. Name three ways to lose data on all of them at once.
The two edges of backups meet here: the backup that saves you from event 29 is the same one that keeps Alice's erased address. Retention has to be long enough to recover from mistakes and short enough (or encrypted per user) to honour deletion.
What to remember from Part 9
- Repair must read a copy known to be at or past what it repairs, and a delete is a versioned write whose tombstone outlives any absence.
- Backups make deleted data restorable until they expire.
- Replication is not a backup: it copies mistakes as fast as it copies writes.
Part 10. Choosing
Designs often pick a replication family by database brand: "we use Cassandra, so we're highly available", or "global tables, so every Region can write". The real choice is who may write a key and what an OK must survive. This Part puts the families and the modes side by side on equal terms.
The same write, three families
The v2 write, "9 Elm St", with the tiny example's timings (the client's own hop to the leader or coordinator not counted):
| Family | How v2 is accepted | OK after | What a read right after may see |
|---|---|---|---|
| Single-leader, quorum (Raft W = 2 of 3) | A appends and sends to B and C; commits at 2 of 3 | max(1.0, 1.4) = 1.4 ms | The leader: v2. A follower: v2 only once applied (C: not until 400 ms) |
| Multi-leader, asynchronous | The local leader commits; other leaders receive it later | 1.0 ms (local commit) | The local leader: v2. Another leader: the old value until replication arrives |
| Leaderless, W = 2 of N = 3 | The coordinator sends to A, B and C; OK at the second answer | 2nd fastest of 1.0, 1.4, 6.0 = 1.4 ms | R = 2: v2 (overlap). R = 1: maybe the old value |
Families at equal durability
Comparing "async multi-leader" with "sync single-leader" compares durability levels, not families. So each family gets an asynchronous row and a synchronous or quorum row:
| Family and mode | Write latency | Writes during a partition | Conflicts | Strongest read | Failover | Repair | Also needs |
|---|---|---|---|---|---|---|---|
| Single-leader, async | Leader's local commit | Only the side with the leader; if a new leader is promoted, the old one must be fenced | None by construction | Leader read (linearizable only with read index or lease) | Promotion loses the unshipped tail | Followers replay the log or a snapshot | Fencing, lag monitoring |
| Single-leader, quorum | max(local sync, (W − 1)-th fastest follower); PostgreSQL: local flush + a round trip | Majority side only | None by construction | Read index or lease on the leader | Election; nothing committed is lost | Same | Elections, fencing |
| Multi-leader, async | Local commit | Every side | LWW, one writer per key, or merge types | Local leader only; stale for keys written elsewhere | Local: nothing to do; the other Region's unreplicated writes are missing until it returns | Leaders exchange changes | A conflict rule |
| Multi-leader, synchronous (MRSC-like) | A round trip to another leader | Only where another leader (a majority of the sites) is reachable | Refused at write time (conflict errors, retried) | Latest version anywhere | Nothing acknowledged is lost | Built in | Three sites (or two plus a witness) |
| Leaderless, W = 1 or small | Fastest replica | Every side | Versions: LWW or siblings | None strong | No failover: any replica serves | Read repair, anti-entropy, tombstone grace | Versions, repair |
| Leaderless, quorum | W-th fastest replica | Only where W home replicas are reachable (strict) | Versions: LWW or siblings | Overlap, not linearizable (Part 5) | Same | Same | Same |
The bytes are the same. Replication sends every write to N − 1 other copies in every family and every mode; a synchronous mode waits for those bytes, it doesn't send more of them. The hook's 20 MB/s of cross-AZ traffic and its $1,036.80 a month are the same whether the two replicas are synchronous or not. The families differ only at the edges: the client's hop to the leader or coordinator (a cross-AZ hop for clients in other AZs), and reads with R > 1, which fetch from several copies (digest reads shrink that).
Sync, semi-sync and async on equal terms
| Synchronous (quorum or named standby) | Semi-synchronous (MySQL) | Asynchronous | |
|---|---|---|---|
| Commit latency | A follower round trip; in parallel with the local sync in Raft-style logs, after it in PostgreSQL | A replica round trip, waited for after the binary log sync (whether events leave before the sync is not documented) | Local only |
| Followers down | A quorum carries on; a fixed standby blocks writes | Carries on with any replica; with none, falls back to async after the timeout | No effect on writes |
| Loss when the leader fails | Nothing acknowledged | Nothing acknowledged, unless it had fallen back | The unshipped tail (Formula 1) |
| Freshness of replica reads | Not fresh: stored is not applied (except PostgreSQL remote_apply) | Not fresh | Not fresh |
| Bytes | Same | Same | Same |
| Fits | Data whose OK must survive a failover: money, orders, ownership | Same, when a replica outage must not become a write outage | Caches, analytics, reads that tolerate loss; cross-Region copies where a round trip per write is too slow |
Conflicts when many can write
With more than one writer, two writes to the same key can both succeed, and something must decide:
| Option | How it decides | Cost |
|---|---|---|
| Last writer wins (LWW) | The write with the latest timestamp wins everywhere. DynamoDB global tables in the default mode use "the latest internal timestamp on a per-item basis" | The other write disappears with no error |
| One writer per key | Each key has a home (a Region or a leader); only the home writes it, and the home moves with an epoch (page 03) | Writes to a key pay the trip to its home; the handover must be fenced |
| Merge types (CRDTs) | Data types whose concurrent updates always merge the same way (counters, sets); the key-value loop (step 2.2) keeps siblings with vector clocks instead | Only works for data that can be modelled that way; an application design, not a database setting |
Two more traps with asynchronous multi-leader, both from the loops' fact-checks: a condition is checked only against the local copy (DynamoDB: "Conditional writes evaluate the condition expression against the version of the item in the Region"), so two Regions can both pass the same check; and a transaction is atomic only where it ran (DynamoDB: "only atomic within the Region where the operation was invoked"), so another Region can briefly see part of it.
Alice in two Regions
The two-Region replay forks from the state after event 8. addr:alice is an item in a DynamoDB global table in the default mode (multi-Region eventual consistency, MREC), with replicas in Regions R1 and R2; both hold "9 Elm St" at version 8. Each device writes with an optimistic check, "only if the version is still 8". We assume a replication delay of 800 ms (AWS: "typically within a second or less"). From the reference implementation:
| # | g (ms) | Event | R1 | R2 |
|---|---|---|---|---|
| G1 | g0 | Laptop, in R1: write "5 Ash Ct" if version = 8. Passes on R1's copy | "5 Ash Ct" v9 | "9 Elm St" v8 |
| G2 | g300 | Phone, in R2, within the replication delay: write "7 Birch Ave" if version = 8. Passes on R2's copy | "5 Ash Ct" v9 | "7 Birch Ave" v9 |
| G3 | g800 | The laptop's write reaches R2: its timestamp (g0) is older than R2's (g300): discarded | "7 Birch Ave" | |
| g1,100 | The phone's write reaches R1: newer: replaces "5 Ash Ct" | "7 Birch Ave" | "7 Birch Ave" |
Both conditional writes succeeded, both clients were told OK, and the laptop's address vanished without an error. Under multi-Region strong consistency (MRSC) (with a third Region or a witness), the laptop's write is replicated to at least one other Region before it returns, and "Conditional writes always evaluate the condition expression against the latest version of an item": the phone's check at g300 sees version 9 and fails, so the phone re-reads and the user decides. (A write to an item that is still being modified in another Region fails with ReplicatedWriteConflictException, which can be retried.)
Why not make every service multi-leader across Regions, and resolve conflicts later?
For what happens when a whole Region fails (detection, routing, RPO and RTO), see the Multi-Region Failover loop primitive (coming); drill 05, "The Booking That Existed in Frankfurt but Not in Virginia", belongs to that page and is further reading after the replay above.
What to remember from Part 10
- Pick the family by who may write a key, not by the database brand.
- Sync and async send the same bytes; sync pays latency and availability to wait for them.
- Multi-leader needs a conflict rule: last writer wins, one writer per key, or data types that merge.
Part 11. On AWS
Every AWS database uses replication; what differs is where it says OK, what a read can see, and what a failover can lose. Stated only from AWS's public documentation and papers.
Managed services that use it
| Service | What it replicates, and when it says OK | Reads | What AWS states |
|---|---|---|---|
| Amazon Aurora | Redo to one storage volume with six copies across three AZs; a write needs four of them; the read quorum is used only for recovery (WAL page, Part 11). Aurora Replicas share the volume, so promoting one loses nothing committed | Up to 15 Aurora Replicas; lag "usually much less than 100 milliseconds", higher under heavy writes (AuroraReplicaLag). The reader endpoint is a stale-read path | Aurora Global Database: one primary Region and up to 10 read-only secondary Regions, replicated "with latency typically under a second"; switchover (planned) with no data loss vs failover (unplanned), whose data loss "depends on the Aurora global database replication lag". Aurora PostgreSQL's rds.global_db_rpo sets an upper bound on RPO by pausing commits on the primary. Write forwarding lets a reader forward writes to the writer, with a read-consistency setting per session (aurora_replica_read_consistency for MySQL, apg_write_forward.consistency_mode for PostgreSQL; SESSION gives read-your-writes for the session's forwarded writes) |
| Amazon RDS | Multi-AZ DB instance: a synchronous standby in another AZ, which can't serve reads. Multi-AZ DB cluster: a writer and two readable readers in three AZs, "semisynchronous replication, which requires acknowledgment from at least one reader". Read replicas: "Amazon RDS copies them asynchronously" (the hook's case) | Cluster readers and read replicas are stale by ReplicaLag. Cluster failover picks the reader "which has the most recent change record"; flow control can throttle the writer so lag doesn't grow unbounded | Data transferred between AZs "for replication of Multi-AZ deployments" is free; replication to read replicas in the same Region is also free |
| Amazon DynamoDB | Each partition is a replication group with a leader (Multi-Paxos); a write is acknowledged "once a quorum of peers persists the log record" (USENIX ATC 2022). "Only the leader replica can serve write and strongly consistent read requests" | Eventually consistent by default (half the price); ConsistentRead on tables and LSIs; GSIs and streams are always eventually consistent | Global tables, MREC (default): asynchronous, "typically within a second or less"; last writer wins per item; conditions checked on the local version; transactions atomic only in the invoking Region; strongly consistent reads can be stale for items last written in another Region; RPO "equal to the replication delay between replicas, usually a few seconds"; ReplicationLatency metric. MRSC: synchronous to at least one other Region; exactly three Regions (three replicas, or two plus a witness); RPO 0; no TTL, no LSIs, no transactions; ReplicatedWriteConflictException. Writes replicated from other Regions bypass DAX. Billing: replicated writes are charged in every replica Region; no cross-Region transfer fee for global-table replication; a witness adds no write, storage or transfer cost |
| Amazon ElastiCache (Valkey, Redis OSS) | "Asynchronous replication mechanisms are used to keep the read replicas synchronized with the primary"; up to five replicas per shard | Replica reads are eventually consistent | A promotion can lose acknowledged writes. With durability turned on (Valkey, node-based clusters), synchronous mode persists each write in a Multi-AZ transactional log before replying; asynchronous mode can lose up to 10 seconds of writes (WAL page, Part 11) |
| Amazon MemoryDB | "Successful write operations are durably stored in a distributed Multi-AZ transactional logs before returning to clients" | Primaries are strongly consistent, "preserved across primary failovers"; replicas are eventually consistent, and reads from a single replica are "sequentially consistent" | A clean example of a quorum-durable log with asynchronous read replicas |
| Amazon MSK | Kafka: a leader and in-sync replicas; acks=all waits for the in-sync set, and min.insync.replicas sets its minimum (WAL page, Part 10). MSK's defaults on a 3-AZ cluster: replication factor 3, min.insync.replicas 2 | Consumers see only messages replicated to the in-sync set | unclean.leader.election.enable defaults to true on MSK clusters without tiered storage (Part 7). No charge for replication traffic between brokers |
| Amazon Keyspaces | Cassandra-compatible; writes are replicated three times across AZs and acknowledged at LOCAL_QUORUM | Choose per read: ONE or LOCAL_ONE may be stale; LOCAL_QUORUM sees prior successful writes | Maps to Part 5's W + R > N only; AWS documents no hinted handoff, sloppy quorum or tombstone grace setting for it, so don't assume them |
| Amazon EKS | The managed control plane runs etcd, "three etcd instances across three AWS Availability Zones": a managed Raft cluster you depend on but never tune | Kubernetes API reads; etcd itself serves linearizable reads by default, serializable (possibly stale) on request | The AWS mapping for Parts 6 and 7 |
Running it yourself
| Option | What it is | Sizing and notes |
|---|---|---|
Amazon EC2 with Amazon EBS: PostgreSQL (streaming replication, synchronous_standby_names = 'ANY 1 (…)' across AZs, a failover manager with a consensus store), MySQL semi-sync, Cassandra or ScyllaDB (replication factor 3, one replica per AZ), etcd or ZooKeeper (3 or 5 nodes), or Kafka | The same mechanics; you own the manager, the fencing, the rewind, the monitoring and the cross-AZ bill | etcd's defaults: 100 ms heartbeat, 1,000 ms election timeout, and an election timeout of at least 10 times the round trip. PostgreSQL's synchronous commit is local flush plus a round trip. MySQL semi-sync falls back after 10,000 ms by default |
| Sizing in words | Every replica needs the leader's write throughput in disk and apply CPU; single-threaded apply sets a lag floor under bursts. The leader's egress carries every write N − 1 times plus any snapshot. Cross-AZ bytes = write rate × size × (N − 1), at 1,037 a month. Synchronous replication adds a round trip per commit, not bytes | Throttle snapshot, rebuild and repair streams; put replicas on the same volume class as the leader; alarm on shipping lag, apply lag and a heartbeat row |
Look-alikes that are not this mechanism
| Look-alike | Why it looks like replication | Why it isn't |
|---|---|---|
| EBS volume replication | "EBS is replicated" | Copies inside one AZ behind one volume: you can't read them or fail over to them, and an AZ loss takes them all |
| EBS snapshots, AWS Backup, RDS automated backups | "Copies of the data" | Point-in-time backups: the protection against event 29, and the reason erased data stays restorable. Not live replicas |
| Amazon S3 strong read-after-write consistency | "S3 is strongly consistent, so its replication is too" | That promise is within one bucket. S3 Replication (cross-Region and same-Region) is asynchronous |
| DAX, or ElastiCache used as a read-through cache | "A copy of the table" | A cache: it doesn't see replicated writes, and serves what it holds until its TTL |
| Route 53 health checks, ARC routing controls, Global Accelerator | "Failover" | They move traffic, not data; the old writer can still write if anything reaches it |
| AWS DMS, zero-ETL integrations, DynamoDB Streams | "They copy changes" | Change streams for consumers (page 05), not replicas that can take over |
| The Aurora volume's six copies | "Six replicas" | Six storage copies behind one writer instance; you read them through instances |
What to remember from Part 11
- Know each service's promise: Aurora and RDS copies, DynamoDB reads and global-table modes, ElastiCache's async replicas and durable mode, MSK's unclean-election default.
- MREC checks conditions locally and resolves by last writer wins; MRSC is three Regions, no TTL, no transactions.
- Replicated writes bypass DAX; a reader endpoint is a replica, not the writer.
Part 12. What you've learned
Back to the 500 lost orders and the stale checkout
Our orders database had two asynchronous replicas 100 ms behind: 500 acknowledged orders at risk in any failover, and checkout pages reading old addresses. Here is how the pieces fix each, and what they cost:
- The 500 orders. Asynchronous replication says OK before any other copy has the write (Part 2). A quorum (commit on 2 of 3) plus an election rule that only lets an up-to-date copy win (Part 6) means no acknowledged write is lost to one failure: in our trace, A died with v3 half-sent, and C, the copy that had it, won. The cost: 1.4 ms per commit instead of 1.0 in a Raft-style log, or 2.4 ms in PostgreSQL, which ships after its local flush; and writes stop if two of three copies are gone. The cross-AZ bytes, 20 MB/s and about $1,037 a month on EC2, were being paid already.
- The stale checkout. C had v2 on disk 4.4 ms after the commit but didn't apply it for 298.6 ms (Part 3, snapshot S1). A position token in the write's OK, checked by the router, sent the read to B, which had applied it (Part 4). Sticky routing would only have made it consistently wrong.
- Around those two:
W + R > Ngives overlap, not linearizability, and a sloppy quorum loses even that (Part 5, S7); a partitioned leader can append but never commit, and its entries are deleted when it rejoins (Part 7, S4 and S5); a late follower catches up from the log or a throttled snapshot (Part 8, S6); repair must read a copy past the delete, tombstones must outlive any absence, and backups keep what replication erased (Part 9, S8); and last-writer-wins across Regions silently drops one of two acknowledged writes (Part 10).
The whole story, event by event
Main line (t in ms), from the reference implementation:
| # | t (ms) | Event | Result |
|---|---|---|---|
| 1 | 0.0 | Laptop writes v1 "12 Oak St" (1.1) | OK at 1.0 (async), 1.4 (Raft W = 2), 6.0 (W = 3), 2.4 (PostgreSQL synchronous); MySQL not drawn |
| 1x | 1.2 | A dies, per mode | Async: an acknowledged write lost. Raft: B has v1 at 1.2, C at 5.8: kept without an OK. PostgreSQL: B has it at 2.2: kept |
| 1y | side | B powered off | Raft 6.0; PostgreSQL ANY 1 (B, C) 7.0; a list naming only B waits; MySQL commits on C's acknowledgement |
| 1z | side | B and C down | Raft: client timeout; PostgreSQL: waits; MySQL: async after 10,000 ms |
| 2 | 50.0 | Report query blocks C's apply until 400 | |
| 3 | 100.0 | Laptop writes v2 "9 Elm St" (1.2) | Committed and OK at 101.4; on C's disk at 105.8 |
| 4 | 110.0 | Heartbeat carries commit 2 | B applies v2 at 110.2 |
| 5 | 150.0 | Read on C | "12 Oak St": stale by 48.6 ms |
| 6 | 180.0 | Read on B | "9 Elm St" |
| 7 | 200.0 | Read on C | "12 Oak St": monotonic reads broken |
| 8 | 400.0 | C applies v2 | Receive lag 4.4 ms; apply lag peaked at 298.6 ms |
| 9 | 990.0 | Last heartbeat that reaches both followers | Received at 990.2 |
| 10 | 1,000.0 | Laptop writes v3 "3 Pine Rd" (1.3); the append to B is lost | On C's disk at 1,005.8 |
| 11 | 1,003.0 | A loses power | v3 on A and C only, not committed |
| 12 | 1,150.2 | B's timer: term 2, vote synced, request sent 1,151.2 | C refuses (more up-to-date), after persisting term 2: refusal at 1,157.2 |
| 13 | 1,240.2 | C's timer: term 3, vote synced 1,245.8 | B grants; C leads term 3 at 1,247.2 |
| 14 | 1,247.2 | C appends no-op 3.4; B refuses (no entry 3) at 1,247.4; C resends 1.3 and 3.4 | Index 4 committed at 1,252.8, index 3 with it; write gap 249.8 ms |
| 15 | 1,500.0 | A restarts | Follows term 3; 3.4 on its disk at 1,501.2 |
| 16 | 2,000.0 | C's heartbeat acknowledged by A and B; C's AZ cut off at 2,001.0 | C's lease ends 2,140.0 |
| 17 | 2,010.0 | Tablet writes v5 "1 Cedar Ln" to C (3.5) | On C's disk at 2,015.6; never committed |
| 18 | 2,140.0 | C's lease ends | 10.2 ms before the earliest possible timeout (2,150.2) |
| 19 | 2,170.2 | A's timer: term 4 | A leads at 2,172.6; no-op 4.5 committed 2,174.0 |
| 20 | 2,200.0 | Phone writes v6 "22 Maple Dr" (4.6) | OK at 2,201.4 |
| 21 | 2,250.0 | Reads at C | No check: "3 Pine Rd" (stale); read index: fails at 3,000.2; lease: refused |
| 22 | 3,000.0 | Link restored | C steps down and deletes 3.5 at 3,000.2; 4.5 and 4.6 on its disk at 3,005.8 |
| 23 | 4,000.0 | B's host fails | B's log ends at 6 |
| 24 | 4,100 to 8,000 | Other writes, indexes 7 to 14 | Each commits in 6.0 ms (C's disk) |
| 25 | 6,000.0 | A snapshots at 10 | A's log starts at 11 |
| 26 | 9,000.0 | B returns needing 7 | Snapshot@10 sent 9,000.4, installed 9,801.6; 11 to 14 stored 9,803.0 |
| 27 | 10,000.0 | Report blocks C's apply; DELETE addr:alice (15) | Committed 10,001.4; B applies 10,010.2; C shows "22 Maple Dr" until 10,300.0 |
| 28 | side | Last night's snapshot and the PITR window | Still hold "22 Maple Dr" |
| 29 | 12,000.0 | Migration sets every address to "" (16) | Committed 12,001.4; all three wrong by 12,010.2 |
Leaderless replay (q ms): Q1 write q3 with C's copy lost, OK at q1.4 · Q2 R = 1 on C stale; R = 2 returns q3 and repairs C · Q3 C down at q15 · Q4 strict write fails, sloppy write OK at q21.2 via stand-in D · Q5 R = 2 on B and C returns q3 · Q6 hints delivered at q100 · Q7 tombstone q5 at q206.0, hint lost at q300, tombstones purged at about q1,201 and q1,206, Merkle repair at q1,400 copies q4 back (with a 2,000 q-ms grace, the tombstone wins).
Two-Region replay (g ms): G1 R1 writes "5 Ash Ct" if version 8 · G2 R2 writes "7 Birch Ave" if version 8 at g300 · G3 by g1,100 both Regions hold "7 Birch Ave"; the laptop's OK'd write is gone.
The cheat card
| Topic | Remember |
|---|---|
| Positions | Stored (on disk) ≥ applied; committed is decided by the leader; reads see applied |
| Async loss | Writes at risk = commit rate × shipping lag (Formula 1) |
| Commit latency | Raft-style: max(local sync, (W − 1)-th fastest follower). PostgreSQL: local flush + standby round trip |
| Followers down | Quorum carries on; a fixed standby blocks; MySQL semi-sync falls back to async after 10 s |
| Lag | Shipping lag decides failover loss; apply lag decides stale reads; detect with a heartbeat row |
| Read-your-writes | Return the position with the OK; read only from a copy applied at or past it, else wait briefly, else the leader |
| Sticky routing | Monotonic reads, not read-your-writes |
| Overlap | W + R > N over the same N home replicas (Formula 2); not linearizable |
| Election | Terms; one vote per term, on disk before sending; only an up-to-date log wins |
| Commit rule | Count copies only for current-term entries; a new leader commits a no-op |
| Partition | Minority appends, never commits; conflicting entries deleted on rejoin |
| Leader reads | Read index (a round trip) or a lease from the heartbeat's send time, ending before any election can finish |
| Replica count | 3 tolerates 1, 5 tolerates 2, 4 is no better than 3 |
| Catch-up | Log while it reaches back; else snapshot at s, then the log from s + 1; throttle the snapshot |
| Deletes | Tombstones outlive any absence plus one repair cycle; repairs read a copy past the delete; backups keep deleted data |
| Not a backup | Replication copies mistakes in milliseconds; keep PITR or a delayed replica |
| Multi-leader | LWW drops a write silently; conditions and transactions are local; decide with one writer per key or MRSC |
Failure checklist
- Does every OK mean what the product promises: which copies, synchronous or not, and what happens when followers are down?
- Is a semi-sync fallback to asynchronous alarmed?
- Are cancelled or timed-out commits treated as unknown outcomes, with idempotent retries?
- Are shipping lag and apply lag both measured, plus a heartbeat row that catches a replica that stopped receiving?
- Does every read that must see the user's own write carry a position, and does every decision read the leader?
- Is there a cache in front of replicas that skips the position check (including DAX with global tables)?
- Are quorums strict where the design relies on overlap, and is
W + R > Ncounted over home replicas? - Does failover fence the old primary before promoting, and pick the most caught-up copy?
- Is an unclean election the Kafka version of promoting a replica that lacks committed entries? Yes: set
unclean.leader.election.enable=falsefor data you promised to keep (MSK defaults it to true without tiered storage). - Are snapshot, rebuild and repair transfers throttled below the leader's spare network and disk?
- Is the tombstone grace period longer than any replica's possible absence plus one repair cycle, and do repairs read a copy past the delete?
- Do backups and PITR windows match both the recovery need and the deletion promise?
Think-first drills
Drill 1. N = 5, W = 3, R = 2. Is a read guaranteed to meet the latest acknowledged write? Which W and R would make it so, and what is the write latency if the coordinator's own replica syncs in 1.0 ms and the other four answer at 1.2, 1.5, 4.0 and 9.0 ms?
Drill 2. A five-node Raft cluster; S1, the leader of term 3, crashes. The logs (term.index): S1: 1.1 1.2 3.3 3.4 3.5 · S2: 1.1 1.2 3.3 3.4 · S3: 1.1 1.2 3.3 · S4: 1.1 1.2 2.3 2.4 2.5 2.6 · S5: 1.1 1.2 3.3. Who can win the next election, and which entries will be deleted?
Drill 3. WAL reaches an asynchronous replica 300 ms late at 2,000 commits a second; the replica is used for reads and as the failover target. How many acknowledged writes can a failover lose, how would you guarantee read-your-writes on it, and what does each change cost?
Interview questions
| Question | Model answer |
|---|---|
What does W + R > N guarantee, and what doesn't it? | That every read set shares at least one of the N home replicas with every acknowledged write's set, so in normal operation a read meets the latest acknowledged write. It isn't linearizability: concurrent writes are settled by versions or last-writer-wins, a failed partial write may or may not be read, a read-then-write is not compare-and-set, and a sloppy quorum counts stand-ins, which breaks the overlap until hints are delivered. |
| Walk me through a Raft election after a leader crash, and why no committed entry is lost. | A follower whose randomized timeout passes increments its term, votes for itself (term and vote on disk first) and asks for votes. A voter grants one vote per term, persists it, and refuses any candidate whose log is less up-to-date (later last term, or same term and not shorter). A committed entry is on a majority, a winner needs a majority, and the two share a node, which refuses a candidate that lacks the entry. The new leader appends a no-op in its term; when that commits, earlier entries commit with it. |
| A user saves a profile and the next page shows the old one. Why, and how do you fix it without sending every read to the primary? | The write went to the primary and the read to a replica that hadn't applied it yet: apply lag, not a bug in either. Return the commit position with the write, keep it in the session (server-side for cross-device), and route the read to a replica whose applied position is at or past it, with a short bounded wait and a fallback to the primary. Sticky routing only gives monotonic reads. Decisions read the primary. |
| Synchronous vs asynchronous replication: what does each lose, and when is each right? | Async acknowledges after the local commit: a failover loses the unshipped tail (commit rate × shipping lag). Sync or quorum acknowledges once enough copies have it: nothing acknowledged is lost to one failure, but each commit pays a round trip (after the local flush in PostgreSQL), a fixed standby's outage blocks writes, and semi-sync may fall back to async. Both send the same bytes. Use sync for data whose OK must survive a failover; async for caches, analytics and far Regions. |
| The network splits a 5-node cluster 2/3. What happens on each side, and when it heals? | The 2-node side can't commit: a commit needs 3 of 5. An old leader there keeps appending entries that time out, and without a read index or lease check it can serve stale reads. The 3-node side elects a leader in a higher term and carries on. On healing, the old leader sees the higher term and steps down; its entries that conflict with the new leader's log (same index, different term) are deleted. Nothing is merged. |
| Why isn't replication a backup? | Replication copies every write, including a bad migration, an accidental delete or a poison entry, to every copy within milliseconds. It protects against independent failures of machines and zones, not against mistakes or correlated failures. Recovery needs a copy from before the mistake: point-in-time restore or a delayed replica. And the backup that saves you also keeps data you deleted, until it expires. |
Where to go next
- Write-Ahead Log, fsync & Group Commit: what one replica's "on disk" means, the acknowledgement table and the write half of the layer trace (Part 10).
- Leases, Fencing Tokens & Distributed Locks: fencing whatever an old leader touched outside the log, clock margins for leases, and lock services under failure (Part 8).
- Idempotency & Effectively-Once Processing: retries for writes whose outcome is unknown.
- Change Streams & the Transactional Outbox: lagging consumers, snapshot-plus-stream bootstraps and repair sources (Parts 7 and 8).
- The Sharding, Hot Keys & Rebalancing and Multi-Region Failover loop primitives (coming): replication within one shard is this page; Region-level detection, routing, RPO and RTO are that one.
- Primitive #09: Consensus, Raft and Paxos, for background, and Primitive #21: Isolation levels, for what session guarantees are not.
- Drill: The Network Partition That Elected Two Leaders, answered in Part 7.
- Loops: the key-value store (steps 1.6, 2.1 to 2.6, R3.8), the message queue (steps 1.3, R1.8, R1.11, R2.7), the digital wallet (steps R2.3, 2.4, 3.4), payments (steps R2.5, 3.5), the stock exchange matching engine (steps 1.4, 2.3, 2.4), the URL shortener (steps 3.1, 3.2), the news feed (steps R1.9, 3.1, 3.2), the paging library (steps R2.3, R2.8) and the Netflix case study (steps 1.3, R1.9, 3.4).