skip to content

Consistency Models

The spectrum of guarantees between strong and eventual: linearizability, causal and session guarantees, convergence, quorums, conflict resolution, CRDTs and reconciliation. Interviewers use this area to find out whether eventually consistent means anything precise to you.

part ofDistributed & scalable systemsoverview, primer and where to startread it →
on this pageshow

questions

page 1 of 2

In a distributed key-value store, what does 'causal consistency' guarantee about the order in which different replicas observe writes, and how does that differ from plain eventual consistency?

level: juniorimportance: must knowfreq 65%

answer

  1. happens-before relation
  2. cause visible before effect, always
  3. concurrent writes = no order guarantee
  4. session consistency in Cosmos DB
  5. dependency tracking cost

basics

~20 s

Causal consistency guarantees that writes which are causally related (one happened because of, or after seeing, another) are seen by every replica in that same order. Unrelated writes may appear in different orders on different replicas. Plain eventual consistency promises no ordering at all, just eventual agreement.

solid answer

~40 s

Causal consistency is a model between eventual and strong (linearizable) consistency. If write A 'happens-before' write B — because the same client wrote A then B, or a client read A then wrote B, or transitively through a chain of such relationships — every replica must apply A before B. Writes with no such relationship (concurrent writes) can be observed in different orders on different replicas, converging eventually via a conflict-resolution rule. Plain eventual consistency only guarantees replicas converge to the same final state if writes stop; it says nothing about the order a client sees writes arrive in, so a reply could be visible before the comment it replies to. Causal consistency rules that out cheaply, without the global coordination strong consistency needs.

go deeper

for a junior

Should state the basic guarantee in plain language: causes are seen before effects, and give one concrete anomaly example (like a reply appearing before its post) that eventual consistency alone permits.

for a middle

Should articulate the happens-before relation precisely (program order, read-from, transitivity) and be able to say what happens to concurrent writes under causal consistency.

for a senior

Should discuss the cost side: dependency tracking overhead, causal stalling, and be able to compare causal consistency's guarantees against a concrete system (e.g., session consistency in Cosmos DB) and explain when it's insufficient.

for a principal

Should place causal consistency in the broader consistency spectrum and articulate the deeper reason it's a sweet spot (availability trade-offs), and know how to choose it vs. stronger/weaker models at a system-design level for a given product requirement.

## What the model guarantees **Causal consistency** is a consistency model for replicated data stores that sits strictly between **eventual consistency** (the weakest useful guarantee) and **strong/linearizable consistency** (the strongest, and most expensive). Its core idea is the `happens-before` relation, borrowed from Lamport's work on distributed event ordering. An operation A happens-before an operation B if: - **(1) Program order** — the same client performed A and then B in program order. - **(2) Read-from** — B is a read that observed the value written by A (a read-from relationship). - **(3) Transitive** — the relation is transitive: A happens-before C if A happens-before B and B happens-before C. Two operations related by happens-before are **causally related**; two operations related in neither direction are **concurrent**. Causal consistency's guarantee: every replica must apply causally related writes in the order happens-before dictates — no replica ever observes an effect before its cause. Concurrent writes carry no such obligation — different replicas can apply them in different orders, as long as they eventually converge to the same value via a deterministic conflict-resolution rule (e.g., `last-writer-wins` or an application merge function). ## Why the model exists This model exists because plain eventual consistency is too weak to build usable applications on. Its only promise is: if writes stop arriving, every replica eventually holds the same data. It says nothing about the order a client sees writes arrive in along the way. That opens the door to real anomalies: - a user could load a comment thread and see a reply ('Great point!') rendered above the comment it replies to, because the reply happened to propagate to that replica first; - a chat client could show 'yes' before the question it answers. These are exactly the class of bug unstructured eventually-consistent systems produce when they replicate independently per key. Causal consistency closes this gap cheaply: it only tracks and respects dependencies that are actually causally linked, not a single global order on all operations, avoiding the expensive coordination (consensus, locking, a single serialization point) strong consistency needs. ## The trade-off The trade-off is real on both sides. - In exchange for a far more intuitive **programming model** — 'you never see an effect before its cause' — the system gives up any agreed ordering for concurrent writes. - Two users editing the same document at once from different regions can each see their own edit applied first, then reconciled later; an application that cannot tolerate that (e.g., a bank ledger where every debit must be globally ordered against every credit) cannot rely on causal consistency alone. - On the systems side it isn't free either: preserving happens-before requires tracking, for every write, which prior writes it causally depends on, and withholding a write from visibility on a replica until all its dependencies are already visible there — that bookkeeping and enforcement both cost metadata and sometimes latency. ## The failure mode in production The most common production failure mode is a symptom of that bookkeeping: a system that cannot yet prove a write's dependencies have arrived at a given replica must either delay exposing the write there (increasing tail latency, sometimes called **causal stalling**) or risk violating the guarantee. Under network partitions or straggler replicas, dependency chains can pile up and cause visible writes to lag behind what a naive eventually-consistent system would have shown immediately — engineers sometimes mistake this for a bug rather than the cost of the guarantee. ## Where it shows up - **Azure Cosmos DB's 'Session' consistency level** is a well-known concrete example, a practical, client-scoped variant of causal consistency: within a single client session, reads are guaranteed to reflect that session's own prior writes and move monotonically forward, at far lower latency and availability cost than the database's 'Strong' consistency level. - Academic systems like **COPS** demonstrated a full, cluster-wide causal consistency model could be built with modest overhead while remaining available during partitions, which is why causal consistency is often cited as the strongest model achievable without sacrificing availability.

  • If causal consistency is stronger than eventual consistency, why doesn't every distributed database just use it by default?
    Because tracking causal dependencies and enforcing 'dependency before visibility' costs metadata and can add latency when a replica must wait for a write's dependencies to arrive before exposing it. For workloads that genuinely don't care about ordering (e.g., independent counters or unrelated keys), plain eventual consistency is cheaper and sufficient, so many systems default to it and offer causal consistency as an opt-in stronger mode.
  • How does a client typically observe a causal-consistency violation in practice?
    A classic symptom is a 'reply before post' or 'answer before question' anomaly — a client reads a value derived from an earlier write without ever seeing that earlier write, because the two writes propagated to that replica out of order. Another is a client re-reading its own prior write and getting a stale value, which violates the read-your-writes session guarantee causal consistency is meant to provide.

Like a group text thread where replies always show up after the message they're replying to, but two people's unrelated jokes posted at the same moment might land in a different order for different readers — nobody is confused about who's replying to what, but simultaneous unrelated chatter can interleave differently.

saying these in an interview costs you the question

  • Says causal consistency means all replicas apply every write in the same global order (that's strong/sequential consistency, not causal).
  • Thinks causal consistency guarantees linearizable reads.
  • Believes eventual consistency already implies causal ordering.
  • Cannot explain what 'happens-before' means or give an example of two concurrent operations.

context

open as a page

In an eventually consistent data store, two replicas each accept a write to the same key while a network partition separates them. When the partition heals and the store uses last-write-wins (LWW) to reconcile, what happens to the value, and what has to happen for that to work reliably?

level: juniorimportance: must knowfreq 75%

basics

~20 s

Last-write-wins keeps whichever write has the newer timestamp and throws away the other one. For this to work, every replica needs a timestamp on each write, and those timestamps must be trustworthy enough to compare across machines.

open as a page

In a distributed database, what does 'eventual consistency' mean, and how is it different from a system that guarantees you always read the latest write?

level: juniorimportance: must knowfreq 85%

basics

~10 s

Eventual consistency means all copies of data will match up eventually, but right after a write some copies might still be stale. Strong consistency means everyone sees the newest value immediately.

open as a page

In an eventually consistent distributed database, what does it mean for replicas to 'converge' after a write, and what is one basic mechanism that helps them get there?

level: juniorimportance: must knowfreq 70%

basics

~20 s

Convergence means all copies of the data eventually end up matching each other, even if some copies see the update late. One simple way this happens is copies periodically comparing notes and fixing whichever one has stale data.

open as a page

What problem does a CRDT (Conflict-free Replicated Data Type) solve, and how does a G-Counter's merge rule guarantee replicas converge without any coordination between nodes?

level: juniorimportance: must knowfreq 70%

basics

~20 s

A CRDT lets many computers update the same piece of data at the same time without talking to each other, and their copies can always be combined back into one correct answer later. A G-Counter only grows, so merging two copies means keeping the bigger count each machine reported.

open as a page

In plain terms, what does it mean for an operation on a single data item to be 'linearizable'?

level: juniorimportance: must knowfreq 70%

basics

~10 s

Every read or write on that item behaves as if it happened instantly, at one moment, and once a write is confirmed, nobody sees an older value again.

open as a page

In a replicated data store, what is a read/write quorum, and what does the condition R + W > N guarantee (where N is the number of replicas, W is the write quorum size, and R is the read quorum size)?

level: juniorimportance: must knowfreq 80%

basics

~20 s

A quorum is the minimum number of copies of your data that must respond before a read or write counts as 'done'. If the number of copies you write to plus the number you read from is bigger than the total number of copies, you're guaranteed to always read at least one copy that has the newest write.

open as a page

You post a photo on a social media app, then immediately refresh the page and the photo is missing. What session guarantee is being violated here, and what would need to change for the app to satisfy it?

level: juniorimportance: must knowfreq 70%

basics

~10 s

This breaks 'read-your-writes' — after you write something, you should always be able to read it back yourself, even if other users still see stale data for a while.

open as a page

In a causally consistent data store, how does a write's set of causal dependencies get determined, and what rule must the store enforce before it lets a replica expose that write to readers?

level: middleimportance: must knowfreq 55%

basics

~20 s

A write's dependencies are every earlier write the client has read or written before making this write. The store must make sure a replica shows all of those earlier writes first, before it shows the new one — otherwise readers could see an effect without its cause.

open as a page

A user posts a comment from their phone, immediately refreshes the page, and the refreshed view does not show their own comment even though other users' unrelated comments do load. Which causal-consistency guarantee is being violated here, and what does 'session causality' mean in general?

level: middleimportance: must knowfreq 60%

basics

~20 s

This breaks 'read your own writes' — a client should always see its own prior actions reflected back to it. Session causality bundles that together with a few other same-client guarantees (like never seeing time move backward for yourself) so a single user's experience is always self-consistent, even if the wider system is only eventually consistent.

open as a page

A distributed key-value store attaches a vector clock to every value it stores, one counter per replica. When two versions of the same key are compared during a read, how does the store use the vector clocks to decide whether one version happened-before the other, or whether the two are concurrent and represent a real conflict?

level: middleimportance: must knowfreq 70%

basics

~20 s

A vector clock is a small list of counters, one per replica, attached to a value. By comparing the lists between two versions, the system can tell if one grew directly out of the other, or if they happened independently and truly conflict.

open as a page

A user posts a comment on a social app, then immediately refreshes the page and sees the comment has vanished. Which consistency guarantee is missing, and how would a system provide it on top of an eventually consistent backend?

level: middleimportance: must knowfreq 70%

basics

~20 s

The app doesn't guarantee 'read-your-writes' - showing your own updates back to you right after you make them. Fixes include reading from the same replica you wrote to, or checking a version marker before answering a read.

open as a page

In a peer-to-peer store that uses gossip to propagate data (not to detect failures), how does the protocol actually spread a write and repair a replica that missed it?

level: middleimportance: must knowfreq 55%

basics

~20 s

Nodes randomly pick another node every so often and compare/share their data. If one has something the other is missing, they copy it over. Repeating this many times across the cluster spreads every update everywhere without needing a central coordinator.

open as a page

In a quorum-based eventually consistent store, what triggers read-repair versus hinted handoff, and how do the two mechanisms differ in how they help replicas converge?

level: middleimportance: must knowfreq 60%

basics

~20 s

Hinted handoff kicks in when a replica is down during a write — another node holds onto the write and delivers it once the replica comes back. Read-repair kicks in during a read — if the replicas queried disagree, the coordinator fixes the stale one on the spot.

open as a page

What is the structural difference between a state-based CRDT (CvRDT) and an operation-based CRDT (CmRDT), and what does each approach require from the network layer to guarantee convergence?

level: middleimportance: must knowfreq 60%

basics

~20 s

One kind of CRDT sends its whole current value to other computers, and they combine values together (state-based). The other kind sends just 'what changed' as a small message, applied in a way that works regardless of order (operation-based). The first tolerates lost or duplicate messages better; the second sends less data but needs more careful delivery.

open as a page

What is a compare-and-set (CAS) operation, and why is it typically the primitive used to build linearizable writes and distributed coordination on top of a replicated store?

level: middleimportance: must knowfreq 65%

basics

~20 s

Compare-and-set writes a new value only if the current value still matches what you expected. It lets many clients race to update the same item safely, because only one 'wins' at a time -- no lost updates.

open as a page

How does linearizability differ from serializability, and how can a database be serializable but not linearizable, or linearizable but not serializable?

level: middleimportance: must knowfreq 85%

basics

~20 s

Serializability makes whole transactions touching many pieces of data behave as if run one at a time, but ignores real-world timing. Linearizability keeps one piece of data instantly up to date in real time, but says nothing about multi-item transactions.

open as a page

Consider a key replicated on 5 nodes (N=5) with write quorum W=3 and read quorum R=2. Walk through why a read is, or is not, guaranteed to return the most recently acknowledged write. How would the answer change if R were increased to 3?

level: middleimportance: must knowfreq 70%

basics

~20 s

With 5 copies, writing to 3 and reading from only 2 isn't enough — 3+2 just equals 5, so the 2 you read from could be exactly the 2 that were skipped by the write. Bump the read to 3 copies and now 3+3=6 is more than 5, so the two groups must share at least one copy.

open as a page

In an eventually consistent, multi-replica data store, two replicas need to figure out which keys have diverged after being offline from each other, without transferring their entire datasets. Explain how a Merkle tree is used to do this efficiently.

level: middleimportance: must knowfreq 55%

basics

~20 s

Each replica builds a tree of hashes over its data — leaves hash small chunks of data, and each parent hashes its children up to one root hash. Two replicas compare hashes top-down; matching hashes mean that whole branch is identical, so they only dig into branches whose hashes differ.

open as a page

A distributed key-value store is eventually consistent. Explain what each of these client-side session guarantees promises — monotonic reads, read-your-writes, monotonic writes, and writes-follow-reads — and give one concrete bug that appears when each is absent.

level: middleimportance: must knowfreq 65%

basics

~20 s

They're four separate promises to one user's session so nothing feels broken: reads won't go backward in time, you'll always see your own writes, your writes won't get reordered, and anything you read before writing will be visible to anyone who later sees that write.

open as a page

When automatic conflict resolution (like LWW or vector-clock causality checks) cannot safely pick a winner because two writes are truly concurrent and semantically different, systems like Amazon's Dynamo return both versions to the application as 'siblings.' Concretely, how does an application perform a semantic merge on such siblings, using a shopping-cart example, and what does the application have to guarantee for the merge to be safe?

level: seniorimportance: must knowfreq 65%

basics

~20 s

When a database can't safely decide which of two conflicting writes is 'right,' it hands both versions to the application, and the application's own code — which understands what the data means — merges them into one sensible result, like combining two shopping carts instead of picking one and losing items from the other.

open as a page

Last-write-wins is popular because it's simple and requires no application involvement. Walk through at least three concrete ways LWW can silently lose committed data or violate causality in production, and describe one situation where you would deliberately choose LWW anyway despite these risks.

level: seniorimportance: must knowfreq 60%

basics

~20 s

Last-write-wins can quietly throw away real changes when clocks disagree, when a whole record gets overwritten instead of just the changed part, or when a 'delete' loses to an older 'update.' It's still fine to use when only the newest value ever matters, like a cache or a heartbeat.

open as a page

In a distributed comment system, causal consistency is used so that a reply always appears after the comment it's replying to, even though the two might be written to different replicas. What mechanism enforces this ordering, and what does causal consistency NOT guarantee?

level: seniorimportance: must knowfreq 55%

basics

~20 s

Causal consistency makes sure that if one event depends on another, like a reply on a comment, everyone sees them in that order. It doesn't force unrelated events from different users to show up in the same order for everyone.

open as a page

What three algebraic properties must a CRDT's merge function satisfy to guarantee replicas converge, and how does an LWW-Register (last-write-wins register) satisfy them while still having a well-known failure mode?

level: seniorimportance: must knowfreq 50%

basics

~20 s

A CRDT's merge rule has to give the same answer no matter which order you combine copies in (commutative), no matter how you group multiple merges together (associative), and merging something with itself again changes nothing (idempotent). An LWW-Register just keeps whichever write has the latest timestamp - that math works, but it means one of two concurrent writes is always thrown away, which can silently lose data.

open as a page

In an OR-Set (observed-remove set) CRDT, how are add and remove operations tracked so that a concurrent add and remove of the same element resolves correctly, and why does a naive two-phase-set (an add-set plus a tombstone remove-set) fail on this case?

level: seniorimportance: must knowfreq 55%

basics

~30 s

An OR-Set tags every 'add' with a unique ID, and a 'remove' only cancels the specific adds it has actually seen. So if one replica adds an item while another replica, not knowing about that add yet, tries to remove the same-named item, the new add survives because its unique tag was never observed by that remove. A simpler design that just marks a name as 'removed forever' would wrongly block that later add too.

open as a page

What latency and availability costs does a distributed system pay to offer linearizable reads and writes, and when would you deliberately choose not to pay them?

level: seniorimportance: must knowfreq 78%

basics

~20 s

Linearizability needs every operation checked against a single up-to-date source of truth -- a leader or majority of replicas -- and that round trip adds delay. If unreachable, the system must refuse requests rather than risk a wrong answer.

open as a page

When tuning the read quorum size R and write quorum size W for a fixed replica count N in a quorum-replicated store, what are the concrete trade-offs between read latency, write latency, and consistency strength? Contrast R=1,W=N with R=N,W=1 and with a balanced choice like R=W=(N/2)+1.

level: seniorimportance: must knowfreq 65%

basics

~20 s

Making reads wait on more copies makes reads slower but writes faster, and vice versa for writes. Waiting on all copies for one side makes that side slow and fragile (one down node blocks it), while the balanced middle spreads the cost evenly across reads and writes.

open as a page

Why is causal consistency often cited as the strongest consistency model that a distributed data store can provide while still remaining available under network partitions, and what does a team give up if they choose a stronger model instead?

level: principalimportance: must knowfreq 35%

basics

~20 s

Because checking 'did this depend on something else' can be answered using only information the client and the local replica already have, without waiting on other machines. Anything stronger — like a single global order for all writes — needs machines to talk to each other and agree in real time, which becomes impossible if a network split cuts them off from each other.

open as a page

A team building a Dynamo-style store advertises 'vector clocks' for conflict detection, but a reviewer says what they've actually implemented is a version vector, and that the distinction matters for correctness. What is the difference between a vector clock and a version vector in this context, and what unbounded-growth problem do both share in a system with many nodes or clients?

level: middleimportance: should knowfreq 40%

basics

~20 s

A 'version vector' is basically the same idea as a vector clock, but it counts per storage replica instead of per client, which keeps it small and stable no matter how many different users write to the data.

open as a page

How does a PN-Counter support both increment and decrement using only grow-only counters internally, and why doesn't it work to just merge a single counter value using max even if deltas can be negative?

level: middleimportance: should knowfreq 55%

basics

~20 s

A PN-Counter is really two separate 'only goes up' counters glued together - one counts all the increments, one counts all the decrements - and the real value is increments minus decrements. You can't just keep one number and take the max across replicas, because max can't tell a smaller number from 'this replica went down on purpose' versus 'this replica hasn't caught up yet.'

open as a page

showing 1–30 of 45