CAP theorem
SystemsHigh priority~30 minNot a menu of three. A conditional: while the network is dropping messages between replicas, refuse some requests or answer with data you cannot prove is current.
Definition
The CAP theorem is a claim about replicated data under network failure. Eric Brewer conjectured it in 2000; Gilbert and Lynch proved it in 2002. A store that keeps more than one copy of the same data cannot hold consistency, availability and partition tolerance at once.
Stated as a menu of three, that is close to useless. Stated as a conditional, it is the sharpest tool in distributed systems: while the network is losing messages between replicas, you must either refuse some requests or answer with data you cannot prove is current.
Both halves of the cost are real. Choosing consistency costs availability, but only inside the partition window — outside it you pay nothing. Choosing availability costs reconciliation: a merge rule, plus compensation for whatever happened while the copies disagreed. That bill comes due at design time, because a system that never planned for divergence cannot un-fork it afterwards.
This is where redundancy sends its invoice. A single copy has no CAP problem at all. The second copy buys survivability and, in the same stroke, makes every replica agrees and every replica answers jointly unsatisfiable the moment a link drops — the concrete case of the trade system design is built on.
When to use
Which one a design wants is usually legible in the prompt. Data carrying an invariant — a balance that must not go negative, a seat that must not be double-booked, a username that must be unique — is asking for C. A feed, a cart, a like count, a profile: staleness is cheap there and A is the default.
Ask it per operation rather than per system: the same store usually runs both, and the question is which of this request's answers you would rather be wrong.
Techniques
The three properties, as the proof states them
Gilbert and Lynch fix all three terms tightly, and the tightness is what makes the theorem say anything at all.
Consistency is atomic, equivalently linearizable: there must exist a total order on all operations such that each looks as if it completed at a single instant. Availability admits no exceptions and no deadline — every request received by a non-failing node must result in a response, with no bound on when. Not usually, not within some target: every request, every reachable node. Partition tolerance means the network may lose arbitrarily many messages between nodes and the guarantees must still hold.
Read that way they are not three interchangeable options. P is not a property you buy; it is a hazard the network hands you. What is left is the conditional — while messages are being lost, you get C or A.
| Consistency | Availability | Partition tolerance | |
|---|---|---|---|
| Where it is configured | Quorum sizes and consistency levels, set per call | Not set directly — it is whatever C leaves behind | Not a setting at all; a fact about the deployment |
| Scope of the guarantee | One key or one transaction at a time | The whole node: it answers, or it does not | The whole cluster, for as long as messages are lost |
| How a violation shows up | A client reads a value another client already overwrote | A reachable node errors, hangs, or never replies | Both sides act as if they were the whole cluster |
| Give it up and you get | AP: answers everywhere, divergence to reconcile | CP: a provably current answer, or none | A design that is only honest on a single node |
PACELC — the branch that runs when nothing is broken
CAP speaks only during a partition, which leaves most of a system's life uncovered. Daniel Abadi's PACELC supplies the missing clause, verbatim:
"if there is a partition (P), how does the system trade off availability and consistency (A and C); else (E), when the system is running normally in the absence of partitions, how does the system trade off latency (L) and consistency (C)?"
His 2012 paper classifies systems on both letters: default Dynamo, Cassandra and Riak as PA/EL; VoltDB/H-Store, Megastore and BigTable/HBase as PC/EC; MongoDB as PA/EC; PNUTS as PC/EL. Products drift, so quote those as the paper's snapshot rather than today's defaults.
Abadi and Brewer are not opponents. Both hold that the classic framing is too coarse, and they differ on which axis it under-serves. Brewer adds what a system does during and after the partition; Abadi adds that replication charges consistency for latency even when nothing is broken — "CAP is only one of the two major reasons that modern DDBSs reduce consistency."
Related concepts
CAP names the axis but not the points on it. Which guarantee an AP store still offers, and what linearizable actually promises, belong to consistency models; the pick one for this workload framing belongs to strong vs eventual consistency. How a partition is detected and bounded is network partitions.
Worked examples
Redis Cluster: the minority side gives up on a timer
Redis Cluster shards across masters, replicates asynchronously, and documents its partition behaviour rather than claiming a letter. Split it, and the majority side promotes a replica of any master it can no longer reach. The minority side keeps accepting writes for up to NODE_TIMEOUT; those writes are acknowledged and then thrown away when the old master rejoins as a replica of the promoted one.
So the same cluster is AP for a bounded window, with acknowledged-write loss, and CP afterwards. The spec also prices the exposure. Take five masters, each with one replica. After one node is lost, just one of the remaining 2N − 1 = 9 nodes is its replica, so a second loss hits it with probability 1/9 — about 11% — and takes a shard down.
Dynamo: always writeable
Amazon's Dynamo takes the other branch by policy. Its shopping-cart workload made a rejected write the worse outcome — the paper's reasoning is that "rejecting customer updates could result in a poor customer experience" — so add to cart has to succeed even when replicas cannot reach each other. Conflicting versions are retained, tagged with vector clocks, and reconciled when something reads them; for a cart the merge rule is union, which can resurrect a removed item but never drops an added one.
Its quorum is a dial rather than a stance: the paper's common configuration is (N, R, W) = (3, 2, 2), so R + W > N and the read and write sets overlap, with both kept under N for latency. Same partition as Redis Cluster, opposite answer, and both are published positions rather than folklore.
Spanner: technically CP, effectively CA
Brewer's 2017 note on Spanner is the useful edge case, and he states the classification himself: "during (some) partitions, Spanner chooses C and forfeits A. It is technically a CP system." What follows is an empirical argument, not a loophole — Google's private network makes that forced choice rare enough that users can treat the system as CA.
The evidence he gives: "there were no events in which a large set of clusters were partitioned from another large set of clusters", the network accounts for 7.6% of Spanner incidents, and Chubby measures 99.99958% availability over outages of 30 seconds or more.
Tradeoffs
| CP — refuse during the partition | AP — answer during the partition | |
|---|---|---|
| What the minority side does | Rejects writes, and usually reads too | Serves and accepts writes locally |
| What the client sees | Errors and timeouts, never a stale committed value | An answer — possibly stale, possibly one of several versions |
| Work owed afterwards | None; state never forked | Reconciliation and compensation |
| Where it bites | An invariant-free workload pays availability for nothing | Every invariant needs a merge rule, and some have none |
| When the cost ends | Once the link heals and a leader is re-elected — the election itself is the CP branch's lingering cost | Later: divergence outlives the partition until something merges it |
Picking A is not picking do nothing. Brewer's 2012 account gives the AP side a three-phase job: detect the start of the partition, enter an explicit partition mode that limits which operations run, and then run a recovery that restores state and makes good the mistakes committed while divergent.
His ATM example is the honest version. Withdrawals are capped at a small amount rather than blocked, and overdrafts that slip through are settled afterwards with a fee and an expectation of repayment. The invariant is not preserved — it is restored, and something outside the system absorbs the gap. An AP design missing that second half is an outage with extra steps.
Things to look out for
- "Pick two of three" — and its sequel, "we're a CA system." You never pick P; the network hands it to you. Brewer junked the phrasing himself in 2012: "the '2 of 3' formulation was always misleading because it tended to oversimplify the tensions among properties." Under the proof's definitions CA means not partition tolerant — undefined behaviour the first time a link drops — so one node is the only honest CA deployment.
- Pinning a letter on a whole database. The call is per operation, and most stores expose the dial: Dynamo's
RandW, Cassandra's consistency levels, MongoDB's read and write concerns. One cluster can run a CP write path beside an AP read path. - Confusing CAP's A with an availability target. The theorem admits no percentile and no measurement window; a 99.99% figure is a fraction of time computed over one. Different claim, different arithmetic, and Availability's subject rather than this page's.
- Reading CAP's C as ACID's C. ACID's consistency says a transaction leaves the declared invariants intact; CAP's is linearizability, a claim about ordering across replicas. A single-node Postgres is fully ACID-consistent and has no CAP C to discuss at all (ACID transactions covers the other one).
- Picturing a partition as a severed cable. Brewer's operational definition is a clock: "a partition is a time bound on communication. Failing to achieve consistency within the time bound implies a partition." To the node waiting on a timeout, slow and dead are one event — and the cause is more often a bad deploy, a firewall rule or an overloaded box than a cut cable.
In the interview
- "Is this design CP or AP?" The strong answer declines that granularity and re-asks per operation: the payment ledger is CP, because during a partition an error beats a double-spend; the catalogue and the view counts are AP, because stale is fine. Then name what the AP half owes — a merge rule, and compensation for what happened while divergent.
- "Spanner is consistent and highly available. Doesn't that break CAP?" No, and the fact is the cheap half of the answer. Spend it on the caveat: effectively CA is bought with a privately owned network and a measured partition rate. None of that transfers to a design running over the public internet.
- "What does CAP say about your design right now, with no partition?" Nothing — and saying so plainly is the strongest move available on this topic. Follow it with PACELC's else-branch: the tradeoff that is live today is consistency against latency, and it shapes far more of the design than the partition branch ever will.