Why 'eventually consistent' became acceptable
For decades, anything weaker than strict consistency was a bug. Then the internet got big enough that strict consistency stopped being affordable — and a generation of engineers learned to live with the gap.
On this page
The picture version
Six pictures for a reader who has never run a distributed system, following one comment: the one you posted that vanished on refresh.
1 · The problem
You post a comment. Refresh. It’s gone.
2 · The naive way
One authoritative box is correct — and unusable.
3 · The trade
Ack the write locally. Tell everyone else afterwards.
4 · The bill
Two edits meet, and something has to decide.
5 · The framing, said correctly
CAP is about partitions, not about steady state.
6 · Keep this card
The whole thing on one index card.
Why it exists
You post a comment, hit refresh, and it’s gone. Refresh again — there it is. Nothing was broken and nothing was lost; you just happened to hit a different replica than the one that took your write, and it hadn’t heard the news yet. That comment is the running example for this post: one small write that some servers know about and others don’t, for a window of time nobody advertised to you.
A 1990s database textbook, asked “is it okay for two reads of the same key to return different values?”, would mostly treat that as a bug to be fixed. Databases were single boxes, transactions were ACID, and the entire point of the system was that callers never had to think about the inside.
Then the web happened. Specifically: shopping carts, social feeds, and comment threads that had to stay up across multiple data centers, often on different continents, while serving millions of concurrent users. The old answer — “one authoritative box, everyone talks to it, transactions serialize” — stopped working at that scale, for two stubborn reasons:
- Latency. A coordinated write across continents is bounded below by the speed of light — call it tens of milliseconds each way between Virginia and Frankfurt, and the exact number is a property of the route, not something to memorize. If your “add to cart” needs a quorum across both, every click inherits that floor and no amount of engineering removes it.
- Partitions. Wide-area links fail. Switches reboot, fibers get cut by anchors, BGP misroutes. If your design says “we cannot serve writes during a partition,” then a transatlantic hiccup takes the site down.
Eventual consistency is the deliberate trade: replicas are allowed to disagree for a while, but they’re guaranteed to converge to the same state once writes stop. You give up “every read sees the latest write” in exchange for staying available, fast, and tolerant of the network being the network.
Why it matters now
You almost certainly depend on something eventually consistent today, even if your primary database is Postgres:
- DNS — a classic eventually-consistent system, and the one whose convergence knob is public: the TTL on a record sets how long caches may keep an answer, which bounds how long a stale one circulates.
- Object storage. S3 was famously eventually consistent for list-after-write until December 2020, when AWS announced strong read-after-write consistency for GET, PUT and LIST across all regions. That’s worth knowing precisely, because a lot of defensive code written before then is still in the wild working around a problem that’s gone.
- Everything with a replica. CDN cache invalidation, DynamoDB’s default read mode, Cassandra, search indexes, and every read replica behind every primary database — including yours, if you’ve ever added one to take load off Postgres.
Knowing when “eventually consistent” is fine and when it’s a bug is a core skill for anyone shipping something distributed, which is most software.
The short answer
eventual consistency = replicas may disagree now + guaranteed to converge later
Picture to keep: a rumour spreading through an office. Everyone will know it by lunch, but for the next few minutes the answer to “what happened?” depends entirely on which desk you ask at. Like a rumour, except rumours have no rule for what happens when two incompatible versions meet — and that rule is the part of eventual consistency that decides whether anyone’s data survives.
You’re trading “all readers see the same value at the same instant” for “all readers will agree once writes settle.” Inside that gap lives a lot of practical engineering: how long is the window, what does the application show during it, and what happens when two writes race?
How it works
Build it from the failure, one constraint at a time. Follow your comment.
Naive attempt: one authoritative database, every replica reads from it. Correct by construction — nobody can see a stale comment because there’s only one copy. Why it breaks: the two reasons above. Readers in Frankfurt eat a transatlantic round trip on every page load, and the moment the link flaps, they see nothing at all.
Fix: put a full replica in each region and let it answer reads locally. Why it breaks: now the write has to get there. If you wait for every replica to confirm before acking your comment, you’ve re-imported the latency you were trying to escape, and you’re back to being unavailable during a partition — you just moved the coordination from reads to writes.
Fix: ack the write locally, propagate it in the background. Your comment
is durable on the replica that took it, you get your 200 OK without waiting
for anyone across an ocean, and the update spreads afterwards by gossip, a
replication log, or an
anti-entropy
job that periodically compares replicas and ships whatever’s missing. This
is the whole trade, and everything below is repair work on it. The refresh
that showed no comment hit a replica the propagation hadn’t reached yet.
Why it breaks: two replicas can now accept writes to the same key without knowing about each other — you edit your comment on your phone while a moderator edits it on a laptop. Fix: pick a conflict-resolution rule. This is where “eventual” stops being a question about delay and becomes a choice about your data. There are three families, and the choice is genuinely consequential:
- LWW — compare timestamps, keep the later one. Simplest, and it silently deletes one of the two edits.
- Vector clocks — don’t resolve it, detect it, and hand both versions to the application to merge. Honest, but now your application has a merge UI to build.
- CRDTs — restrict the data to shapes (counters, sets, sequences) where merging is provably order-independent, so there’s nothing to resolve. Automatic, but only for those shapes. See crdts.
The convergence guarantee is real but conditional: it assumes writes eventually stop (or at least that the rate of new writes stays below the rate at which replicas can sync). In practice the window is usually milliseconds to seconds; under partition it can be minutes or hours.
The CAP framing — useful, often misquoted
Most engineers have heard of CAP. The accurate statement, from Eric Brewer’s original conjecture (later proved by Gilbert and Lynch in 2002): during a network partition, a distributed system can either remain available or remain strongly consistent, not both. When the network is healthy, you can have both — which is most of the time, and is exactly why the theorem is so often misapplied to steady-state behaviour it says nothing about.
The pop-culture version — “pick two of CAP” — is misleading. Partitions aren’t optional, so the real choice is what to do when one happens. Systems that pick A under partition (Dynamo, Cassandra, most caches) are the ones we call eventually consistent. Systems that pick C (Spanner, etcd, ZooKeeper) will refuse writes — or refuse reads — rather than diverge.
The PACELC refinement
CAP only describes partition behavior. Daniel Abadi’s PACELC adds the missing case: even when there’s no partition, you still trade latency against consistency. A strongly consistent write that needs a quorum across regions is fundamentally slower than a local-only write. So even healthy systems make this choice every time you tune replication factor or read consistency level.
Show the seams
- “Eventual” is not a duration. The definition says “if writes stop, replicas converge,” and says nothing about when. In a misconfigured cluster, “eventually” can be never. Cassandra is explicit about this in its own docs: hinted handoff is best-effort, and it’s anti-entropy repair that actually guarantees convergence. Skip the repairs and divergent replicas simply stay divergent.
- Read-your-own-writes is not free. A user who just posted a comment and doesn’t see it on refresh thinks the site is broken. Most systems bolt on session consistency via sticky routing or per-session read tokens to paper over this.
- Eventual consistency does not mean ‘no consistency.’ Modern systems layer guarantees on top: monotonic reads, monotonic writes, causal consistency. Causal consistency in particular (“if A happened before B, no replica shows B without A”) is my pick for the best value on the menu — it rules out the anomalies users actually notice, like a reply appearing above the comment it answers, without paying for global coordination. That’s a judgment about typical product needs, not a theorem.
- Strong consistency is available at global scale, by several different routes. Spanner and CockroachDB both offer it across regions, but they do not share a mechanism, and conflating them is a common error. Spanner leans on TrueTime, backed by atomic clocks and GPS, to bound clock uncertainty; CockroachDB explicitly states it does not require special clock hardware and uses a different scheme. What they share is the bill: consensus rounds, and the latency those rounds cost.
- Nobody publishes the split. There’s no measurement of what fraction of production traffic runs on eventually consistent stores versus strongly consistent ones, so this post describes the landscape rather than sizing it. The shape of the trade — latency versus coordination — hasn’t changed in a decade.
You started with eventual consistency = replicas may disagree now + guaranteed to converge later. What did this post add? — + a conflict rule you are choosing whether or not you know it. “Converge later” sounds like a
timing property, so people tune replication lag and stop there. But two
replicas that accepted different writes have to decide something, and where
the default is last-writer-wins — as it is in a number of widely used
stores — the disagreement is resolved by throwing one of the writes away.
Look up which rule your store uses; it is rarely on the marketing page. The
disappearing comment was the harmless version of this; two conflicting edits to
the same comment are the expensive one. Eventual consistency wasn’t a
giving-up — it was an honest acknowledgment that the network has physical
limits. The part worth staying suspicious of isn’t the “eventually.” It’s what
the system quietly did to your data to get there.
Check yourself
Before you go — your team fixes the disappearing-comment bug by pinning each user’s session to one replica. A week later, support reports that two users in the same thread see the replies in different orders. Did the fix cause this?
Answer
No — the fix worked, it just solved a smaller problem than anyone assumed. Sticky sessions buy read-your-own-writes: you always talk to the replica that took your write, so your own comment can’t vanish. That says nothing about other people’s writes, which still arrive at different replicas in different orders. Guaranteeing that a reply never appears before the comment it answers is a different and stronger property — causal consistency — and it needs the system to track happens-before relationships (vector clocks or similar), not just route you to a fixed server. Worth internalising: “eventually consistent” is not one guarantee but a floor, and the useful engineering is in which specific guarantee you bolt on top.
And one more — a colleague argues that since network partitions are rare, they can build a globally-replicated system that is both strongly consistent and always available, and just accept “rare downtime.” Is that a coherent position?
Answer
It’s a coherent engineering position and an incoherent claim about CAP. CAP doesn’t say you must choose in the abstract; it says that during a partition you cannot have both, and partitions are not optional. So your colleague has in fact chosen C — their system will refuse writes during a partition, which is the “rare downtime” they’re accepting. That’s a completely legitimate design; etcd and Spanner make it deliberately. The incoherent part is calling it “always available.” And PACELC names the bill they haven’t counted: even on a perfectly healthy network, agreeing across regions costs consensus round trips on every write, every day, partition or not.
Famous related terms
- CAP theorem —
CAP ≈ during a partition, pick Available or Consistent— the framing that made eventual consistency intellectually respectable. - CRDT —
CRDT = data structure + commutative merge function— makes “eventual” automatic for specific shapes; part of the conceptual toolkit behind collaborative editing. - Strong consistency / linearizability —
linearizable ≈ "behaves as if there's one copy"— the gold standard you’re explicitly walking away from. - Read-your-writes —
RYW ≈ "at minimum, a session sees its own edits"— the practical floor most user-facing eventually-consistent systems implement. - Quorum —
quorum ≈ R + W > N— the classic dial for sliding along the consistency-latency spectrum without changing engines. Treat the inequality as a starting heuristic: sloppy quorums, repair lag, and the conflict rule underneath all weaken what it actually guarantees.
Going deeper
- Gilbert & Lynch, “Brewer’s conjecture and the feasibility of consistent, available, partition-tolerant web services” (2002) — answers “what does CAP actually prove?”, which turns out to be much narrower and much more precise than the version that circulates.
- Werner Vogels, “Eventually Consistent” (ACM Queue, 2008) — answers “how does a company that runs one of these reason about it?”, and is where the vocabulary of session, monotonic-read, and read-your-writes consistency became common currency.
- Daniel Abadi, “Consistency Tradeoffs in Modern Distributed Database System Design” (IEEE Computer, 2012) — the rabbit hole, answering “what does this cost me on a healthy network?” via PACELC.