CAP Theorem
CAP Theorem
Definition: The CAP Theorem (Brewer’s conjecture, formally proven by Gilbert and Lynch in 2002) states that a distributed data store can simultaneously provide at most two of three guarantees once a network partition occurs: Consistency, Availability, and Partition Tolerance.
How It Works
- Consistency (C): every read receives the most recent write or an error. This is linearizability, a real-time recency guarantee, stronger than the “C” in ACID.
- Availability (A): every request to a non-failing node receives a non-error response within bounded time. Nothing in this guarantee says the response is current, only that the node answers.
- Partition Tolerance (P): the system keeps operating despite arbitrary message loss or delay between nodes. TCP/IP networks drop and delay packets routinely; partitions are a certainty at scale, not an edge case.
- Because P is not optional on any real network, the actual design decision is what happens during a partition: CP or AP.
- A “CA” system exists only on a single node, or inside a network that provably never partitions. Both cases disappear the moment data is distributed across machines connected by a fallible network.
- CP systems (etcd, ZooKeeper, HBase, MongoDB with majority read/write concern): when a partition splits the cluster, the minority side stops answering reads and writes until it can confirm it isn’t stale. It sacrifices availability to preserve correctness.
- AP systems (Cassandra, DynamoDB, Riak, CouchDB): when a partition hits, every reachable node keeps serving reads and writes. Replicas can diverge during the partition and are reconciled afterward.
- Reconciliation strategies for AP systems include last-write-wins timestamps (simple, but silently drops the losing write), vector clocks (detect conflicts explicitly, push resolution to the application), and CRDTs (data structures that merge deterministically without conflict, at the cost of restricting what operations are allowed).
- Gilbert and Lynch’s proof uses an asynchronous network model: given a network where messages between two nodes can be delayed indefinitely, no algorithm can guarantee both a linearizable read and a non-error response.
- The core argument is a two-node case: partition the pair, write to one, read from the other. The read either returns stale data (violates C) or blocks/errors (violates A) — there is no third option under an indefinite message delay.
- Quorum-based systems make the choice tunable rather than binary. With N replicas, a write requiring W acknowledgments, and a read requiring R acknowledgments, setting
W + R > Nguarantees every read overlaps at least one up-to-date replica, pushing the system toward CP behavior. Lowering W and R toward 1 pushes it toward AP. - Detecting a partition at all is nontrivial: a node can’t cleanly distinguish “the network dropped my messages” from “the remote node crashed” from “the remote node is just slow.” Most systems use timeouts and heartbeats as a proxy for partition detection, so the CP/AP decision is really being made against a guess.
- CAP is closely related to, but distinct from, the FLP impossibility result (Fischer, Lynch, Paterson, 1985), which shows that deterministic consensus is impossible in a fully asynchronous system if even one node can fail. CAP is about the C/A trade-off under partition; FLP is about consensus termination under failure. Both push distributed systems toward using timeouts, leases, or partial synchrony assumptions to make progress at all.
- “Partition” in CAP means any communication failure between nodes that should be able to talk to each other, not just a full network split. A single dropped or delayed message between two replicas counts.
Trade-offs
- Choosing CP turns partition-time latency into outright unavailability: clients see errors or timeouts, not wrong answers. This fits domains like financial ledgers or inventory counts where a wrong answer is worse than no answer.
- Choosing AP means clients always get an answer, but it may be stale, or in multi-writer systems, may later need to be merged with a conflicting write. This fits domains like social feeds or shopping carts, where staleness is tolerable and blocking is not.
- The trade-off is asymmetric in practice. Partitions are rare, so most of the real cost of a CAP choice is paid every day as a latency tax during normal operation, not during the rare partition window — this is exactly the gap PACELC fills.
- Neither choice removes complexity, it only relocates it. CP pushes complexity to the client, which must handle errors, retries, and unavailability windows.
- AP pushes complexity into the application layer, which must handle conflict detection and resolution: merge functions, CRDTs, or manual reconciliation logic that a CP system would never need.
- Partial CP is possible at the operation level: a system can be CP for writes to a specific record via per-key quorum, while remaining broadly available for reads of unrelated keys. CAP is often discussed as a whole-system property, but production systems frequently apply it per data path.
Why It Matters
- It forces an explicit decision about degraded-mode behavior. Without naming the choice, a system’s partition-time behavior is undefined, and gets discovered during a production incident instead of a design review.
- It shapes what guarantees a client library can safely assume. A CP store lets client code assume read-after-write consistency for free; an AP store pushes conflict handling into application code.
- It directly informs achievable SLAs. AP systems can advertise very high availability because they never block on cross-node coordination to answer a request; CP systems trade some of that away for stronger correctness.
- It gives teams shared vocabulary during incident review. “We went unavailable because we’re CP and lost quorum” is a precise, defensible statement; “the database broke” is not.
- It surfaces early in interviews and design docs because it’s the smallest model that captures why “just use a distributed database” is not a complete answer to a system design problem.
Common Pitfalls
- Applying CAP reasoning outside of a partition. During normal operation the real trade-off is latency versus consistency, which is what PACELC Theorem was created to describe.
- Conflating CAP’s “Consistency” (linearizability, a recency guarantee) with ACID’s “Consistency” (the database enforces its own declared constraints). They share a word and nothing else.
- Treating a database as permanently “CP” or “AP” as a fixed label. Many systems are tunable per query via quorum settings, so the honest answer to “is it CP or AP” is often “depends on the consistency level requested for this call.”
- Assuming AP means no consistency guarantees at all. Most AP systems still offer eventual consistency, and many offer tunable levels that approach strong consistency at the cost of latency.
- Defaulting to CP for a service where a moment of unavailability is worse than serving slightly stale data (a product listing page, a follower count), out of theoretical purity rather than actual business impact.
- Forgetting CAP is scoped to behavior during a partition specifically. It says nothing about steady-state throughput, disaster recovery, or backup strategy — those are separate design concerns entirely.
- Treating “we chose AP” as the end of the design conversation. AP systems still need an explicit conflict resolution strategy; skipping that step just means conflicts get resolved arbitrarily, by whichever write happened to land last.
- Assuming the theorem says anything about latency. CAP is silent on latency by design; that gap is exactly why PACELC exists as a follow-on.
Comparison
| CAP Theorem | PACELC Theorem | ACID Consistency | |
|---|---|---|---|
| Scope | Behavior only during a network partition | Behavior during a partition AND during normal operation | Data satisfies application-defined constraints |
| Trade-off named | Consistency vs Availability | Consistency vs Availability (partition), Consistency vs Latency (else) | N/A, a correctness property, not a trade-off |
| Covers steady-state latency cost | No | Yes | No |
| Typical use | Reasoning about failure-mode design | Reasoning about everyday latency/consistency cost | Reasoning about transaction correctness |
| System | Partition behavior | Classification |
|---|---|---|
| etcd, ZooKeeper | Minority side stops serving | CP |
| Cassandra, Riak | All reachable nodes keep serving | AP |
| DynamoDB (default) | All reachable nodes keep serving | AP |
| MongoDB (majority concern) | Minority side stops serving | CP |
| Google Spanner | Uses TrueTime + quorum to stay consistent, accepts added latency | CP |
| ACID | BASE | |
|---|---|---|
| Stands for | Atomicity, Consistency, Isolation, Durability | Basically Available, Soft state, Eventually consistent |
| Consistency model | Strong, transactional | Eventual |
| Typical system | Single-node RDBMS, CP-leaning distributed stores | AP-leaning distributed stores (Cassandra, DynamoDB) |
| Trade-off favored | Correctness and isolation over raw availability | Availability and latency over immediate correctness |
| Relation to CAP | The natural transaction model for CP systems | The natural transaction model for AP systems |
History
- Eric Brewer presented the CAP conjecture as a keynote at PODC 2000, drawing on observations from building distributed systems at Inktomi, without a formal proof at the time.
- Seth Gilbert and Nancy Lynch formalized and proved it in 2002 (“Brewer’s Conjecture and the Feasibility of Consistent, Available, Partition-Tolerant Web Services”), giving the theorem its rigorous footing and its common name.
- Brewer revisited the theorem in a 2012 paper, “CAP Twelve Years Later,” clarifying that the “2 of 3” framing oversimplifies things: partitions are rare, so the real-world question is how a system behaves during the rare partition window versus how it behaves the rest of the time.
- That clarification is what later got formalized as PACELC, which explicitly separates “during a partition” behavior from “else” (normal operation) behavior.
- The theorem predates most of today’s popular AP and CP databases; Cassandra, DynamoDB, and MongoDB were all designed with explicit awareness of the CAP trade-off, unlike earlier single-master relational systems that never had to reason about it.
Real-World Scenario
A checkout service backed by a CP inventory store loses connectivity between its two data centers. The minority-side replicas stop accepting writes to avoid overselling the last unit of a product. Customers in that data center see “please try again” on checkout instead of a false “in stock” confirmation. Once the network heals, the minority side rejoins, catches up on the replicated log, and resumes serving.
Compare that to a session-store service backed by an AP store: during the same partition, both sides keep accepting writes, a user’s session data briefly diverges between data centers, and it gets merged (or the more recent write wins) once connectivity is restored. Neither service goes fully down, but the checkout service intentionally trades availability for correctness where correctness matters more, while the session store intentionally trades correctness for uptime where a few seconds of stale session state is invisible to the user.
Debugging Walkthrough: Tracing a Partition
A network link between two data centers, us-east and us-west, drops. Nodes on each side stop receiving heartbeats from the other side after a few missed intervals. What happens next depends entirely on the CP/AP choice:
- CP cluster (etcd-backed inventory service, 5 voting members, 3 in
us-east, 2 inus-west):us-westno longer has a majority. Any write attempted there fails fast, returning an error likeetcdserver: request timed outorno leaderwithin the configured election timeout (commonly 1-5 seconds). Reads requiring linearizability are rejected too.us-east, still holding a majority, keeps its leader and serves normally, unaffected. - On-call sees a sharp error-rate spike scoped specifically to
us-west’s write path, a “quorum lost” alert firing for that side, andus-east’s dashboards showing no change at all. The failure is loud, immediate, and geographically isolated to the minority side. - AP cluster (Cassandra-backed session store) hitting the same link failure: both sides keep accepting reads and writes at a local consistency level (
LOCAL_QUORUMorONE). No error-rate spike occurs. Instead, hinted-handoff and hint-replay counters climb on both sides as each accumulates writes the other side hasn’t seen. - On-call sees no availability alert at all during the partition. The only signal is a rise in “read repair” or “hints stored” metrics, and possibly a support ticket about a user seeing a stale or flickering value depending on which side answered their request.
- When the link is restored: the CP cluster’s minority nodes rejoin, catch up on the replicated log from the current leader, and resume serving, typically within seconds of the network healing. The AP cluster keeps serving throughout, but spends extra background CPU and network bandwidth replaying hints and running anti-entropy repair to reconcile the two sides’ divergent writes, a process that can take much longer than the partition itself.
- The postmortem difference is stark: the CP incident has a clear start and end time bounded by the outage window; the AP incident has no downtime to report at all, but a data-reconciliation tail that has to be separately verified as complete.
Design Checklist
When picking a CAP posture for a new data path, work through these questions before touching a database:
- Is a stale read actually harmful here, or just cosmetically wrong?
- Is a failed/blocked request more costly to the business than a wrong answer?
- Can the conflict-resolution logic for an AP choice actually be written correctly, or does it just get deferred to “figure it out later”?
- Does this data path need the same answer as every other path, or can different features on the same product use different postures?
- What does the client do when the answer is “unavailable” — retry, queue, degrade the UI, or fail the whole request?
FAQ
Does CAP mean I have to give up consistency entirely to get availability? No. AP systems typically offer eventual consistency, not zero consistency. The trade-off is about behavior during a partition, not a permanent abandonment of correctness.
Can a system be CP for some operations and AP for others? Yes. Many databases apply the trade-off per query via tunable consistency levels or per-key quorum settings, rather than as one global mode.
Is a single-node database subject to CAP? No. CAP only applies once data is replicated across nodes connected by a network that can partition. A single node has nothing to partition from.
Why do people say “CAP is a false trichotomy”? Because in practice P isn’t a choice, it’s a fact of networked systems. The real, everyday decision is C versus A, and only while a partition is actually happening.
Does choosing AP mean I lose durability too? No, durability (a write surviving a crash) is orthogonal to CAP. An AP system can still fsync to disk before acknowledging a write; CAP is about cross-node behavior during a partition, not single-node durability.
How does client-side caching interact with CAP? It doesn’t change the store’s classification, but it can silently reintroduce staleness into an otherwise CP system if the client serves cached reads instead of round-tripping to the database, so caching policy needs to be reasoned about separately from the store’s own guarantees.
Is Spanner a counterexample to CAP? No. Spanner is CP: it uses synchronized clocks (TrueTime) and quorum writes to stay linearizable, and it does sacrifice some availability and latency during partitions and clock uncertainty windows to do it. It doesn’t escape the theorem, it just engineers around the practical cost of the CP choice.
Common Interview Questions
- What’s the difference between CAP’s Consistency and ACID’s Consistency? A recency guarantee versus a constraint-satisfaction guarantee; they are unrelated despite the shared word.
- Why is a “CA” system considered not to exist in a real distributed deployment? Because any network of more than one node can partition, and CA requires assuming it never will.
- How would you decide between a CP and an AP data store for a given feature? By asking whether a wrong/stale answer or a blocked/failed request is worse for that specific feature’s users.
- What does PACELC add that CAP doesn’t cover? The latency-versus-consistency trade-off that exists during normal, non-partitioned operation.
- How do quorum reads and writes relate to CAP? They make the C/A trade-off tunable per operation rather than fixed at the system level.
- Can you give an example of a system that’s CP for one kind of data and AP for another? Many multi-model databases and most large platforms run a CP metadata/config store (e.g., ZooKeeper/etcd) alongside an AP data store (e.g., Cassandra) for bulk application data.
- Why can’t you just add more replicas to avoid the CAP trade-off? More replicas change the odds and the blast radius, but they don’t remove the fundamental fact that a partition can still separate any subset of them from the rest.
Example
MongoDB configured with readConcern: majority and writes acknowledged by a majority of replica set members behaves as CP: it rejects operations on a minority partition rather than risk serving stale data. Cassandra configured with ONE read/write consistency behaves as AP: every reachable node answers regardless of partition state, and conflicting writes are resolved afterward via last-write-wins timestamps or read-repair.
Related Terms
Referenced by