Strong and eventual consistency.
Copy data to three regions and a question appears that one database never had to answer. When someone reads, which copy's answer counts? Strong consistency makes every copy act like one machine and pays for it in waiting. Eventual consistency answers from the nearest copy at once and lets the copies catch up afterwards. Most real systems use both, choosing per operation.
The idea.
One database never has to decide which copy is right. Three copies do.
One rename, three copies
Tandoor, a recipe-sharing app, keeps every recipe on three replicas: Mumbai, Frankfurt and Virginia. Priya renames her recipe from "Dal Makhani" to "Dal Makhani (no cream)". Twenty milliseconds after her phone shows "Saved", a reader in Virginia opens the recipe. Which name should they see?
Under strong consistency the answer is fixed: the new name, exactly as if there were one copy. Under eventual consistency the reader may still see the old name for a short while, called the inconsistency window. The only promise is that once edits stop, every replica ends up holding the same value. Neither answer is wrong. They are two different contracts, and the difference is paid for in waiting on one side and in surprises on the other.
The same rename, under two promises
Scenario 1 of 2: As described.
- 1 Returns at 122 ms, once a majority has entry 1042.
- 7 Stored, but Virginia does not yet know it is committed.
- 10 Started 20 ms after the save returned, so it must see the rename.
- 13 Without a lease, Mumbai first confirms it is still leader with a heartbeat round to Frankfurt (~120 ms).
Timeline as a list
The same rename, under two promises: 5 lanes, from 0 ms to 360 ms.
- 0–122 ms · Priya (Mumbai) · rename (ok) — Returns at 122 ms, once a majority has entry 1042.
- 2–122 ms · Mumbai (leader) · waiting for a majority
- 2 ms · Mumbai (leader) → Frankfurt, arriving 62 ms: append 1042
- 2 ms · Mumbai (leader) → Virginia, arriving 97 ms
- 62 ms · Frankfurt · has 1042
- 62 ms · Frankfurt → Mumbai (leader), arriving 122 ms: ack
- 97 ms · Virginia · has 1042 — Stored, but Virginia does not yet know it is committed.
- 122 ms · Mumbai (leader) · commit 1042 (ok)
- 142–332 ms · Virginia · confirm commit point
- 142–332 ms · Reader (Virginia) · read → new name (ok) — Started 20 ms after the save returned, so it must see the rename.
- 142 ms · all lanes · Strong: read after save returns
- 142 ms · Virginia → Mumbai (leader), arriving 237 ms: committed index?
- 237 ms · Mumbai (leader) · answers under its lease — Without a lease, Mumbai first confirms it is still leader with a heartbeat round to Frankfurt (~120 ms).
- 237 ms · Mumbai (leader) → Virginia, arriving 332 ms: ≥ 1042
Four names you will meet
How it works.
The two promises come from two different replication paths. Everything between them is a question of which orderings a reader may observe.
How a strong write and read work
Strong rename, then a strong read from Virginia
- Priya's phone → Mumbai replica (leader): rename(recipe 812, …)
- Mumbai replica (leader) → Frankfurt replica: append entry 1042
- Mumbai replica (leader) → Virginia replica: append entry 1042
- Frankfurt replica → Mumbai replica (leader) (reply): ack 1042
- Note over Mumbai replica (leader): 2 of 3 replicas have 1042: commit
- Mumbai replica (leader) → Priya's phone (reply): ok, version 1042
- Reader in Virginia → Virginia replica: get(recipe 812)
- Virginia replica → Mumbai replica (leader): commit index? (read index)
- Note over Mumbai replica (leader): lease held: answers at once (no lease: +~120 ms)
- Mumbai replica (leader) → Virginia replica (reply): commit index ≥ 1042
- Note over Virginia replica: apply up to 1042, then answer
- Virginia replica → Reader in Virginia (reply): "Dal Makhani (no cream)"
The write waits for a majority, not for all three: two of three copies surviving any single failure is what makes the commit durable. The read is the subtle half. Virginia may already hold entry 1042 without knowing it is committed, or may be missing entries it has never heard of. Before it answers a linearizable read it asks the leader for the current commit index and waits until it has applied that far (Raft calls this a read index). The leader has its own doubt: a newer leader may have been elected without it hearing. So before it hands out the index it either confirms it is still leader with a heartbeat round to a majority (Mumbai to Frankfurt, another ~120 ms), or it holds a time-bounded leader lease that rules out a rival and answers at once. Tandoor's numbers assume the lease. Sending every strong read straight to the leader faces the same choice. Each option costs a cross-region round trip or a dependence on bounded clock drift; none is free.
How an eventual write and read work
Eventual rename, read from the nearest copy
- Priya's phone → Mumbai replica (leader): rename(recipe 812, …)
- Mumbai replica (leader) → Priya's phone (reply): ok (local commit only, ~2 ms)
- Mumbai replica (leader) → Frankfurt replica: ship entry 1042 (async)
- Reader in Virginia → Virginia replica: get(recipe 812)
- Virginia replica → Reader in Virginia (reply): "Dal Makhani" (old)
- Mumbai replica (leader) → Virginia replica: ship entry 1042 (async, ~95 ms)
- Note over Virginia replica: all three replicas agree again
The spectrum between them
Each step down the ladder allows one more kind of surprise, in exchange for answering from fewer machines.
| Model | Promise | What it still allows | Example |
|---|---|---|---|
| Linearizable | Each operation appears to take effect at one instant between its call and its return, in real-time order. | Nothing a single copy would not also do. | etcd; Spanner (external consistency, which extends it to transactions) |
| Sequential | All clients see one total order of operations that keeps each client's own order, but not wall-clock order. | A read that starts after another client's write returned can still miss it. | ZooKeeper: writes in one order; a server may serve stale reads unless the client calls sync |
| Causal | Operations that could have influenced each other are seen in that order by everyone. | Two co-authors' concurrent edits appearing in opposite orders to two readers. | MongoDB causally consistent sessions |
| Eventual | If writes stop, all replicas converge on the same value. | A recipe step shown before the ingredient it uses; a value going backwards between two reads. | Cassandra at consistency level ONE; DynamoDB default reads |
Every row implies the rows below it: a linearizable system is also sequential, causal and eventually consistent (Jepsen's hierarchy). The price runs the other way. Jepsen classes linearizable and sequential as impossible to keep available on both sides of a network partition; causal can stay available to a client that keeps talking to the same replica. Session guarantees such as read-your-writes sit between causal and eventual and have their own topic.
Two histories, judged by each model
Scenario 1 of 4: As described.
- 3 Linearizable forbids this: the write returned at 24 ms, before the read began at 34 ms.
Timeline as a list
Two histories, judged by each model: 3 lanes, from 0 ms to 120 ms.
- 0–24 ms · Priya · w(title = new) (ok)
- 24 ms · all lanes · Linearizable: write returned
- 34–62 ms · Meera · r(title) → old (violation) — Linearizable forbids this: the write returned at 24 ms, before the read began at 34 ms.
- 70–94 ms · Sam · r(title) → new (ok)
Converging is a choice, not a given
"Replicas converge" hides a decision. Tandoor's single leader orders every write, so its replicas only ever lag; they never disagree about which write came last. Conflicts appear when replicas accept writes independently, as in multi-leader setups and leaderless stores such as Dynamo and Cassandra. Picture a multi-leader Tandoor: Priya renames the recipe in Mumbai while a co-author renames it in Frankfurt during a partition, both replicas accept a write, and when they reconnect one value must win or both must be kept. Last-writer-wins compares timestamps and keeps the larger one: simple and always convergent, but it silently discards the other write, and with clock skew the discarded one may be the newer. Keeping both versions (siblings) and letting the application merge them loses nothing but pushes work onto the reader: Dynamo merged shopping carts this way, which is why a deleted item could reappear. How versions are tracked (vector clocks) is covered in the key-value store's versioning topic.
One key under eventual consistency
Replicas of one key agree, drift apart while an update travels, and agree again. The conflict branch exists only where more than one replica accepts writes (multi-leader or leaderless); a single leader such as Tandoor's never takes it.
| From → To | Event | Guard | Action |
|---|---|---|---|
| All replicas agree → Update spreading | write at one replica | ||
| Update spreading → Update spreading | another write | same replica | ship in order |
| Update spreading → All replicas agree | last replica applies | ||
| Update spreading → Concurrent versions | write at another replica | multi-leader | |
| Concurrent versions → All replicas agree | replicas sync | last-writer-wins | drop older |
| Concurrent versions → All replicas agree | replicas sync | siblings | app merges |
- All replicas agreestart
- Update spreading
- Some replicas are stale; reads there return the old value.
- Concurrent versions
- Two replicas accepted different writes that neither saw.
In practice.
How often eventual consistency actually bites, what strong costs per request, and the knobs real databases expose.
How often does a stale read actually happen?
- Recipe reads at peak
- 40,000/sassumption
- Recipe edits at peak
- 700/s2M daily users × 10 edits ≈ 230/s on average, × 3 at peak (derived in the read-your-writes topic)
- Distinct recipes read in a day
- 8Massumption
- Replication lag to a replica (p99)
- 200 msassumption
- The editor's reload after saving
- ~100 ms after each editthe app reloads the recipe once the save returns (assumption)
- Recipes mid-replication at any instantwrites × lag = 700 × 0.2 s140from Recipe edits at peak and Replication lag to a replica (p99)
- Chance a random read hits one of theminflight ÷ keys = 140 ÷ 8,000,000≈ 0.0018%from Recipes mid-replication at any instant and Distinct recipes read in a day · Assumes reads spread evenly. Freshly edited recipes are often hotter, so the real figure is higher, but still tiny.
- Editor reloads that can land inside the window (upper bound)writes × 1 reload = 700 × 1 (100 ms < 200 ms lag)≤ 700/sfrom Recipe edits at peak, The editor's reload after saving and Replication lag to a replica (p99) · At p99 lag every reload beats replication. At median lag (tens of ms) most reloads arrive after the window has closed, so the real count is lower.
- Share of all reads that are an editor's rereadown-stale ÷ reads = 700 ÷ 40,000≤ 1.75%from Editor reloads that can land inside the window (upper bound) and Recipe reads at peak
- For strangers, eventual consistency is close to invisible.
- For the person who just wrote, it is likely, and it happens exactly when the user is watching. That is a job for session guarantees (read-your-writes), not a reason to make all 40,000 reads a second strong.
What strong costs per request
- Round trip Mumbai–Frankfurt
- ~120 msassumption
- Round trip Mumbai–Virginia
- ~190 msassumption
- Local durable append
- 2 msassumption
- Local read from the nearest replica
- 1 msassumption
- Strong write (leader plus nearest follower is a majority)disk + rtt-fra = 2 + 120~122 msfrom Local durable append and Round trip Mumbai–Frankfurt
- Eventual write (local commit, ship later)disk~2 msfrom Local durable append
- How much slower a strong write is122 ÷ 2~61×from Strong write (leader plus nearest follower is a majority) and Eventual write (local commit, ship later)
- Strong read in Virginia, leader answering under a leasertt-iad + local = 190 + 1~191 msfrom Round trip Mumbai–Virginia and Local read from the nearest replica
- Strong read in Virginia, no lease (leader confirms leadership with Frankfurt first)rtt-iad + rtt-fra + local = 190 + 120 + 1~311 msfrom Round trip Mumbai–Virginia, Round trip Mumbai–Frankfurt and Local read from the nearest replica
- Eventual read in Virginialocal~1 msfrom Local read from the nearest replica
- Across regions, strong consistency costs a round trip on every write, and on every strong read not served locally by a leader holding a lease; without the lease, the read pays a second round trip.
- The cost is set by geography, not by hardware. Faster disks do not shorten 120 ms of fibre.
Latency per operation, Tandoor's three regions
Data
| Operation | Latency (ms) |
|---|---|
| Eventual read (Virginia) | 1 |
| Eventual write | 2 |
| Strong write | 122 |
| Strong read (Virginia) | 191 |
The knob in real systems
Trade-offs.
The choice is made per operation, not per database. The question is what a stale or reordered read would break.
- Pro:Reads answer from the local region in about 1 ms.
- Pro:Writes keep working in Mumbai even if Frankfurt and Virginia are unreachable.
- Pro:The one visible anomaly (the editor seeing an old version) is fixed by a session guarantee.
- Con:Other readers see an edit up to one replication window late.
- Con:A like count can briefly go backwards if two reads hit different replicas.
About 122 ms on every write and about 191 ms on every strong read from Virginia (311 ms without a leader lease).; Mumbai cannot commit writes if it loses both followers; the eventual design keeps writing as long as Mumbai is up.
- Pro:A credit can never be spent twice, and the balance never goes below zero.
- Pro:Check-then-spend is safe because every read sees every earlier spend.
- Con:Spends take a cross-region round trip.
- Con:A region on the minority side of a partition cannot spend credits until it heals.
Two regions can both approve the last credit; the overspend is found hours later.; Fixing it means refunds, clawbacks and apologies.
Where the promises break in practice
| Failure | Impact | Detection | Mitigation | Meanwhile |
|---|---|---|---|---|
| Leader fails over while asynchronously shipped entries have not reached any follower3Mumbai replica (leader) | Writes that were acknowledged are gone after the new leader takes over. | A gap between the old leader's last log position and the new leader's. | Commit on a majority before acknowledging, or state the recovery point you accept. | |
| Follower falls minutes behind (slow disk, long GC, saturated link)5Virginia replica | Reads there are far staler than the usual 200 ms at p99. | Replica lag metric in seconds and in log positions. | Take the replica out of read rotation above a lag bound. | Virginia reads go to Frankfurt, slower but fresh. |
| Last-writer-wins with clock skew | A newer edit loses to an older one stamped by a machine whose clock runs fast. | Hard; the loss is silent. | Conditional writes on a version number (compare-and-set), or keep siblings and merge. Hybrid logical clocks only stop a fast clock from beating a causally later write; concurrent writes still lose one. | |
| A cache in front of a strong store | The store is linearizable, but what users read is whatever the cache holds. | Cache age of served objects versus the store's version. | Invalidate on write and bound the TTL; see the distributed cache topics. |