Heads up: posts on this site are drafted by Claude and fact-checked by Codex. Both can still get things wrong — read with care and verify anything load-bearing before relying on it.
why → how

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.

Data intermediate Apr 29, 2026 · updated Aug 25, 2026 · 12 min read

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.

refresh refresh you post a comment it’s gone there it is each refresh may land on a different replica — and one hadn’t heard the news yet Nothing was broken, and nothing was lost.
One small write that some servers know about and others don’t, for a window of time nobody advertised to you. That window is what the rest of the post is about.

2 · The naive way

One authoritative box is correct — and unusable.

“one authoritative database; everybody talks to it” the one copy Virginia reader in Frankfurt every page load a transatlantic round trip, every time and when the link flaps × they see nothing at all if the design says “no writes during a partition,” a transatlantic hiccup takes the site down Correct by construction. Unusable at internet scale.
Two stubborn reasons, and neither is a matter of engineering effort: a coordinated write across continents is floored by the speed of light, and wide-area links genuinely fail. Putting a replica in each region only moves the problem to the write path.

3 · The trade

Ack the write locally. Tell everyone else afterwards.

you post 200 OK, no ocean your replica takes the write later, in the background replica in Frankfurt replica in Singapore by gossip, a replication log, or an anti-entropy job This is the whole trade. Everything else is repair work on it. the refresh that showed no comment hit a replica the propagation hadn’t reached yet But two replicas can now take writes to the same key without knowing.
You gave up “every read sees the latest write” and bought back local latency and survival during a partition. The vanished comment is that trade working as designed, not a bug.

4 · The bill

Two edits meet, and something has to decide.

you edit your comment on your phone while a moderator edits it on a laptop your edit their edit concurrent somebody has to decide LWW compare timestamps, keep the later one — and silently delete the other edit vector clocks don’t resolve it, detect it, and hand both versions to the application CRDTs restrict the data to shapes where merging is provably order-independent “Converge later” sounds like a timing property. It is a data policy. and where the default is last-writer-wins, the disagreement is settled by throwing one write away
People tune replication lag and stop there, because “converge later” sounds like it is only about time. But two replicas that accepted different writes have to decide something, and that decision is rarely on the marketing page.

5 · The framing, said correctly

CAP is about partitions, not about steady state.

the accurate statement, not the “pick two” version network healthy consistency and availability, both which is most of the time — and is exactly what the theorem says nothing about during a partition pick one: stay available, or stay strongly consistent Cassandra picks A; etcd and Spanner pick C Partitions aren’t optional, so the real choice is what to do when one happens. PACELC adds the case CAP skips even with no partition you still trade latency against consistency: a quorum across regions is fundamentally slower than a local-only write, so healthy systems make this choice every day too
The pop-culture “pick two of CAP” is misleading because partitions are not something you opt into. Systems that pick availability under partition are the ones we call eventually consistent; the ones that pick consistency refuse writes rather than diverge.

6 · Keep this card

The whole thing on one index card.

eventual consistency = replicas may disagree now + guaranteed to agree later + a conflict rule you are choosing either way ∴ look up which one your store uses the part worth staying suspicious of isn’t the “eventually” — it’s what the system did to your data to get there
Picture to keep: a rumour spreading through an office — everyone will know by lunch, but for the next few minutes the answer to “what happened?” depends on which desk you ask at. Where it breaks: rumours have no rule for what happens when two incompatible versions meet, and that rule is the part which decides whether anyone’s data survives.

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:

  1. 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.
  2. 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:

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:

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

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.

Going deeper