Skip to content
Tech Interview Prep home
Technical interview guide

CAP Theorem & Consistency Models

Why a distributed system can't have perfect consistency, availability, and partition tolerance all at once — and what real systems trade off.

Read
23 min
Practice MCQs
25
Interview QA
25
Edition
v3
Editorial status
Reviewed

Scope: Vendor-neutral distributed-systems principles; current AWS architecture guidance.

Overview

Curated: · Written: · Reviewed:

CAP is a narrow theorem that is usually quoted too broadly

The CAP theorem states that a distributed data store cannot simultaneously guarantee all three of consistency (every read sees the most recent write), availability (every request receives a non-error response), and partition tolerance (the system continues operating despite network partitions).

The framing "pick two" is the source of most misunderstanding. Partition tolerance is not optional. Networks partition — cables fail, switches misconfigure, availability zones become unreachable — and a system that does not tolerate partitions simply stops working when one occurs. So for any system spanning machines, the real choice is what to do during a partition: refuse requests to preserve consistency (CP), or serve requests that may return stale or conflicting data (AP).

Two further clarifications matter more than the theorem itself.

CAP's C is not ACID's C. CAP consistency means every read sees the latest write — closer to linearizability. ACID consistency means a transaction preserves declared constraints. They share a letter and nothing else, and conflating them produces confused conversations.

CAP describes behaviour during a partition only. It says nothing about normal operation, which is the overwhelming majority of the time. This is the theorem's biggest practical limitation, and it is why PACELC is a more useful frame: if there is a Partition, choose between Availability and Consistency; Else, choose between Latency and Consistency. That "else" clause is where most real design decisions live — every synchronous replication choice, every replica read, every cache is a latency-versus-consistency trade made while everything is working normally.

Consistency is a spectrum

Strong consistency (linearizability): every read returns the most recent committed write, as though there were a single copy. Simplest to reason about, most expensive — it requires coordination, which costs latency and reduces availability during partitions.

Sequential and causal consistency sit between. Causal consistency guarantees that causally related operations are seen in order by everyone, while concurrent operations may be seen differently. It is often sufficient and much cheaper than linearizability — if a reply must never appear before the message it replies to, that is a causal requirement, not a linearizable one.

Eventual consistency: if writes stop, replicas converge. The guarantee is about the destination, not the journey, and applications live in the journey.

Between the extremes are the session guarantees, which are frequently what a product actually needs:

  • Read-your-writes: a user always sees their own writes. The most commonly needed and most commonly missing.
  • Monotonic reads: a user never sees data go backwards in time.
  • Monotonic writes: a user's writes are applied in order.

These are much cheaper than global strong consistency and usually eliminate the user-visible symptoms people actually complain about. Sticky routing per session provides most of them almost free.

What eventual consistency does to an application

Read-your-writes fails most visibly. A user updates their profile, the next read hits a lagging replica, they see the old value and conclude the save failed — then repeat it. This is a consistency problem presenting as a usability one, and it is the most common real manifestation.

Monotonic reads fail when consecutive requests hit replicas at different positions, so data appears to go backwards.

Concurrent writes conflict, and something must resolve them. Last-write-wins is simple and silently discards data — acceptable for a presence indicator, unacceptable for a shopping cart. Alternatives are merging, conflict-free replicated data types that merge deterministically, or surfacing the conflict to the user, which is sometimes the only correct answer.

The design responses are mostly cheap: route reads that must be fresh to the leader, or carry a session watermark and avoid replicas that have not caught up; return the written value from the write response rather than re-fetching, which is often the simplest fix and frequently overlooked; choose consistency per operation where the store allows it; and make writes commutative where possible, since an increment merges cleanly and an absolute set does not.

Choosing per operation, not per system

The most important practical point is that consistency is chosen per operation, not once for the whole system. Most stores expose this directly — DynamoDB offers eventually and strongly consistent reads per request; MongoDB has read concerns; quorum systems let you tune read and write quorums.

So the design work is classifying operations. A balance check before a withdrawal needs strong consistency. A follower count on a profile page does not. Charging a card needs it; showing a recommendation does not. The number of operations genuinely needing strong consistency is usually far smaller than a first pass assumes, and paying for it everywhere is the common expensive mistake.

The mirror mistake is assuming eventual consistency everywhere is acceptable because the store defaults to it. A stale read that merely displays something is a cosmetic issue; a stale read that informs a write decision — checking a balance, checking uniqueness, checking inventory — is a correctness bug, and that distinction is the one worth being rigorous about.

Quorums

Many distributed stores use quorums: with N replicas, a write succeeds after W acknowledgements and a read consults R. Under the standard quorum model, R + W > N makes every read set overlap every acknowledged write set. Seeing the latest value also requires reads to compare versions correctly and excludes mechanisms such as sloppy quorums that place replicas outside the intended set. W = N gives fast reads and fragile writes; W = 1, R = 1 gives fast everything and no guarantee.

This makes the trade explicit and tunable rather than a property of the product, which is a useful way to think even about systems that do not expose it directly.

What to take away

CAP is a genuine result and a poor design tool. The useful questions it points at are: what happens during a partition, and is that behaviour deliberate; which operations genuinely need strong consistency, and which merely inherit it out of habit; what does the user see when data is stale, and is that acceptable; and — since these are business tolerances rather than technical ones — how much data loss and how much downtime are actually acceptable, expressed as recovery objectives someone has agreed to.

Worked example: R=1 after W=1 looks like a failed save

N = 3. User writes bio on replica A (W = 1). Next GET hits replica B, 400 ms behind.

readuser seesretries the save?
R = 1, any replicaold bioyes (duplicate write)
sticky session / read-your-writesnew biono
balance check with R = 1 then withdrawstale $50, second debitcorrectness bug

Display staleness is a UX miss. A stale read that informs a write is a CAP mistake, not PACELC folklore.