Consistency Models

How quickly, and how strictly, every copy of your data has to agree with every other copy.

What is it?

Once data is copied across multiple databases (via replication) or split across regions, a natural question comes up: if you write a change to one copy, when — and how reliably — will every other copy reflect that change? Different systems answer this differently, and that answer is called their consistency model.

Strong consistency means every read, everywhere, always sees the latest write — as if there were really only one copy. Eventual consistency relaxes that: a write might take a moment to reach every copy, but given enough time (and no new writes), every copy will eventually agree.

Explain like I'm 10

Strong consistency is like a single shared whiteboard everyone reads from directly — the moment someone writes on it, everyone sees the update immediately. Eventual consistency is like several people taking their own photo of the whiteboard and syncing it later — for a little while, some people's photos might be a version behind, but eventually everyone's photo matches.

Examples

The same read, two different consistency guarantees

// Strong consistency: guaranteed to see the write immediately
await write(userId, { balance: 100 });
const balance1 = await strictRead(userId); // always 100

// Eventual consistency: might briefly see a stale value
await write(userId, { balance: 100 });
const balance2 = await eventualRead(userId); // could still be the old value, briefly

How it works

Strong consistency is usually achieved by routing every read through a single source of truth (or requiring multiple copies to confirm a write before it's considered done), which costs latency and availability during network problems. Eventual consistency instead lets a write settle in on each copy independently, at its own pace, which is faster and more resilient to failures, at the cost of a short window where reads can return stale data.

Why does it exist?

There's a real, unavoidable tradeoff between how fast/reliable a system stays under failures and how strictly up-to-date every read is guaranteed to be — different applications land in different places on that tradeoff depending on how much staleness they can tolerate.

When to use it

Choose strong consistency for data where staleness would cause real harm — an account balance being checked before a withdrawal, an inventory count during checkout. Choose eventual consistency for data where a brief delay doesn't matter much — a like count, a follower count, a search index.

When not to use it

Don't default to strong consistency everywhere out of caution — it costs real latency and availability, and most data in most applications (activity feeds, view counts, recommendations) doesn't actually need it.

Common mistakes

  • Assuming all databases are strongly consistent by default — many popular distributed databases default to eventual consistency for performance.

  • Using eventually-consistent reads for something that genuinely needs strong consistency, like a payment balance.

  • Treating 'eventual' as if it means 'never' — in a healthy system, eventual consistency typically resolves in milliseconds to seconds, not indefinitely.

Practice exercises

  1. Easy:

    Give one example of data where eventual consistency would be perfectly fine, and one where it wouldn't.

  2. Medium:

    Explain, in your own words, why strong consistency tends to cost more latency than eventual consistency.

  3. Hard:

    Describe a real bug that could occur if a signup flow read a user's data eventually-consistently right after writing it.

Interview questions

What's the difference between strong and eventual consistency?

Strong consistency guarantees every read reflects the latest write immediately; eventual consistency allows a brief delay before all copies agree, trading immediacy for speed and resilience.

Why would a system choose eventual consistency over strong consistency?

Eventually consistent systems are typically faster and more available, especially during network issues, since they don't need every copy to confirm before a write is considered complete.

Give an example where eventual consistency would be an acceptable trade-off.

A social media like count or view count — being off by a few for a moment causes no real harm.

What is linearizability, and how is it a stronger guarantee than what people often casually call "strong consistency"?

Linearizability guarantees that every operation appears to take effect atomically at some single instant between when it was called and when it returned, and that all operations across all clients agree on one consistent global ordering — as if there were truly only one copy of the data, updated one operation at a time. It's the strictest, most precisely defined form of strong consistency, not just "reads see recent writes."

What's the difference between strong consistency and causal consistency?

Strong consistency requires every operation to be seen in the same real-time order by everyone. Causal consistency only requires that operations which are causally related (one happened because of, or after seeing, another) are seen in that same order by everyone — but operations with no causal relationship to each other can be seen in different orders by different observers, which is cheaper to guarantee.

Give a concrete scenario where causal consistency matters, but plain eventual consistency alone would produce a confusing result.

A comment reply gets replicated to one reader before the original comment it's replying to does — that reader sees a reply to a comment that, as far as they can tell, doesn't exist yet. Causal consistency prevents this by guaranteeing that if writing the reply causally depended on having seen the original comment, no replica will ever show the reply without also showing the comment it depends on.

What is "monotonic reads" as a consistency guarantee, and what bug does it prevent?

It guarantees that once a client has seen a particular value (or a later one), it will never subsequently see an older value on a later read. It prevents the confusing experience of refreshing a page and having something you'd already seen appear to disappear or revert, because a later read happened to hit a more-lagged replica than an earlier one did.

What is "read-your-writes" consistency, and how is it distinct from monotonic reads?

Read-your-writes specifically guarantees that a client will always see its own prior writes on subsequent reads. Monotonic reads is broader and unrelated to authorship — it guarantees a client never sees time move backward on reads regardless of who made the write, including writes made by other users.

What is "monotonic writes" consistency?

It guarantees that writes from the same client are applied in the order that client issued them, everywhere — so a system never applies a client's second write before its first, which could otherwise leave data in an order the client never intended (e.g. an 'undo' applying before the action it was meant to undo).

What does quorum-based consistency (R + W > N) mean, and why does that formula guarantee an overlap?

N is the number of replicas holding a piece of data, W is how many replicas must acknowledge a write for it to succeed, and R is how many replicas a read must consult. If R + W > N, any set of W replicas that accepted a write and any set of R replicas consulted by a read must share at least one replica in common — guaranteeing every read consults at least one replica that has the latest write.

Scenario: N=3 replicas, W=2, R=1. Is this configuration guaranteed to always read the latest write?

No — R + W = 3, which is not greater than N (3), so there's no guaranteed overlap. A read could consult the one replica that wasn't part of the two that accepted the latest write, returning a stale value.

Scenario: N=3, W=2, R=2. Is this configuration now guaranteed to read the latest write?

Yes — R + W = 4 > N (3), so any 2 replicas consulted by a read must overlap with at least one of the 2 replicas that accepted the write, meaning the read is guaranteed to see it (or a newer value), at the cost of needing to contact 2 replicas per read instead of 1.

What's the difference between "sequential consistency" and "linearizability"?

Both require all clients to agree on one single global ordering of operations. Linearizability additionally requires that ordering to respect real, wall-clock time — if operation A completed before operation B started, A must come first in the agreed order. Sequential consistency drops that real-time requirement: any single global ordering is acceptable as long as everyone agrees on the same one, even if it doesn't match the actual real-time sequence of events.

What is PACELC, and how does it extend the CAP theorem's tradeoff?

PACELC observes that CAP only describes the tradeoff during a network Partition (Availability vs Consistency), but even with no partition at all (Else), a system still has to trade off Latency vs Consistency — waiting for more replicas to agree always costs some latency. PACELC frames consistency as a tradeoff a system faces essentially all the time, not just during failures.

Explain CAP theorem precisely — what does it actually claim?

For a distributed system, during a network partition (some nodes can't communicate with others), you must choose between Consistency (every node sees the same, latest data) and Availability (every request gets a response) — you cannot guarantee both at that moment. It says nothing about normal, non-partitioned operation, where both can often be provided together.

Common misconception: "CAP theorem means a system is either CP or AP, all the time." Why is this an oversimplification?

CAP's tradeoff only actually applies during a partition — outside of one, a well-designed system can be both consistent and available. Many real systems are also tunable per-operation (different requests can ask for different consistency guarantees) rather than being permanently locked into one label, and CAP itself doesn't account for the latency-vs-consistency tradeoff that exists even without a partition (that's what PACELC adds).

What is a vector clock, and what problem does it solve in a system with no single global clock?

A vector clock is a small set of per-replica counters attached to each write, incremented by whichever replica made it. Comparing two vector clocks lets the system determine whether one write causally happened-before another, happened-after, or neither (a true conflict) — solving the problem of ordering events across machines that don't share a perfectly synchronized clock.

What is a CRDT (Conflict-free Replicated Data Type), and how does it let replicas merge conflicting writes without coordination?

A CRDT is a data structure specifically designed so that merging any two independently-updated copies always produces the same, well-defined result, regardless of the order the updates are merged in — for example, a counter that only ever increments can be merged by simply summing each replica's increments. This lets replicas accept writes independently and merge later with no need to coordinate or lock during the write itself.

What is "read repair" in a leaderless replicated system, and how does it help move toward consistency?

When a read consults multiple replicas and notices they disagree, it detects which replica(s) have a stale value and proactively writes the newer value back to them as part of handling that same read. Over time, this quietly nudges lagging replicas back in sync without needing a separate, dedicated repair process.

What is "tunable consistency" (as in systems like Cassandra), and why might different requests in the same app choose different levels?

It lets each individual read or write specify its own required consistency level (e.g. how many replicas must respond) rather than the whole system being locked to one setting. An app might use a strong, higher-quorum read for a payment balance, while using a fast, single-replica read for a low-stakes view counter — matching the guarantee's cost to how much that specific data actually needs it.

Scenario: a shopping cart uses eventual consistency. A user adds an item on their phone, then immediately opens the site on their laptop and doesn't see it. What guarantee is missing, and what would fix this specific flow?

This is a read-your-writes violation — the laptop's read hit a replica that hadn't yet received the phone's write. Fixing it for this user specifically means routing their reads to a replica known to be caught up on their own writes (or to the primary) for a short window, rather than any arbitrary replica.

Scenario: an inventory count uses eventual consistency across regions. Two customers in different regions both see "1 in stock" and both complete checkout. What went wrong, and how would you prevent overselling?

Each region's replica hadn't yet seen the other's near-simultaneous decrement, so both reads (correctly, for their own stale view) saw availability that no longer existed by the time both writes landed — a classic consequence of accepting concurrent writes under eventual consistency. Preventing it requires either a strongly consistent check specifically for the decrement-and-reserve step (routing it through a single source of truth) or accepting the possibility and handling it after the fact (detecting oversells and resolving them, e.g. refunding or backordering one of the two orders).

Why is "eventually consistent" not the same as "eventually correct"? What happens if writes never stop?

"Eventually consistent" only promises that if writes to a piece of data stopped, every replica would converge to the same value given enough time — it says nothing about what that value is being correct in a business sense, and if writes keep arriving continuously, replicas may never actually reach a moment of full agreement, only ever chasing a moving target.

What does "consistency" mean specifically in the CAP theorem, and how does it differ from the "C" in ACID?

CAP's consistency specifically means linearizability — every read reflects the most recent write, as if there were one copy. ACID's consistency means something different: that a transaction only ever moves the database from one valid state (satisfying its constraints and invariants) to another valid one — it's about application-level correctness rules, not about how fast replicas agree with each other.

Mechanically, why can a strongly consistent system become unavailable during a network partition?

Guaranteeing every read sees the latest write typically requires a quorum of replicas to be reachable and in agreement (per the R+W>N logic). If a partition splits the replicas so no side has a majority, the system can't safely confirm a read or write meets that guarantee — so, rather than risk returning stale or conflicting data, it refuses to serve the request at all until the partition heals.

What's the practical latency cost difference between reading from a single node versus reading with a quorum?

A single-node read is one network round trip and returns as soon as that one node answers. A quorum read must contact multiple replicas and wait for enough of them to respond (and often compare their answers), so its latency is bounded by the slowest of those replicas rather than just one — meaningfully slower, especially if replicas span regions with real network distance between them.

What is causal consistency's relationship to "happens-before" ordering?

Causal consistency is defined directly in terms of the happens-before relation: if operation A happens-before operation B (B read a value A wrote, or they occurred in sequence on the same client), every replica must show A before B. Operations with no happens-before relationship to each other are unconstrained and can be seen in different orders on different replicas without violating the guarantee.

Scenario: a chat app shows a reply before the message it was replying to, on one user's screen. What consistency guarantee failed?

Causal consistency — the reply causally depends on the original message (it couldn't have been written without it existing first), so any consistency model weaker than causal (like plain eventual consistency with no causal tracking) can let a replica apply or show them out of that dependent order.

Why do many real-world systems mix consistency levels rather than picking one globally?

Different pieces of data carry very different costs for being briefly stale — a payment balance being wrong is a real, damaging bug, while a view count being off by a few for a second is invisible to anyone. Applying the same, expensive strong-consistency guarantee everywhere wastes latency and availability on data that never needed it; mixing levels spends that cost only where staleness would actually cause harm.

What's a "stale read," precisely, and is it always a bug?

A stale read is a read that returns a value older than the most recently completed write, because it was served by a replica (or cache) that hadn't yet applied that write. Whether it's a bug depends entirely on the data and use case — it's expected and harmless for an eventually-consistent view counter, but a serious correctness bug for something like an account balance used to authorize a withdrawal.

Follow-up: give one scenario where a stale read is completely fine, and one where it's a serious bug.

Fine: a follower count on a profile page being a few seconds behind causes no real harm to anyone. Serious: a fraud-detection check reading a stale "account not yet flagged" status right after it should have been flagged, letting a fraudulent transaction through it was specifically meant to block.

What does it mean for a consistency model to provide guarantees "per session" versus "global" guarantees?

A per-session guarantee (like read-your-writes or monotonic reads) only promises consistency from the point of view of one particular client's own sequence of operations — other clients might still observe a different order relative to each other. A global guarantee (like linearizability) promises one single, agreed ordering that holds for every client simultaneously, which is strictly harder and more expensive to provide.

Why is achieving strict linearizability across multiple geographically distant data centers particularly costly, and what do most geo-distributed systems do instead?

Linearizability effectively requires coordinating with distant replicas on every operation, and the speed of light puts a real floor under how fast that round trip can be across, say, different continents — often tens to low hundreds of milliseconds, added to every operation. Most geo-distributed systems instead accept a weaker model (causal consistency, or region-local strong consistency with asynchronous cross-region replication) for the vast majority of operations, reserving true global linearizability for the specific few operations that truly can't tolerate anything less.