/Interview Study Guide/System design
ConceptsPart of Scalability

CAP theorem

SystemsHigh priority~30 min

Not 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.

ConsistencyAvailabilityPartition tolerance
Where it is configuredQuorum sizes and consistency levels, set per callNot set directly — it is whatever C leaves behindNot a setting at all; a fact about the deployment
Scope of the guaranteeOne key or one transaction at a timeThe whole node: it answers, or it does notThe whole cluster, for as long as messages are lost
How a violation shows upA client reads a value another client already overwroteA reachable node errors, hangs, or never repliesBoth sides act as if they were the whole cluster
Give it up and you getAP: answers everywhere, divergence to reconcileCP: a provably current answer, or noneA design that is only honest on a single node
The three terms, on the axes that actually separate them.

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.

Client AReplica 1Replica 2Client Bwrite x = 2replicate x = 2message lost — the partition starts hereAP: ack now · CP: no ack, the write blocks or failsread xAP: x = 1, stale · CP: error, cannot prove it is currentlink heals — reconcileAP owes a merge rule here; CP has nothing to merge
One write, one partition, two branches. The reconcile step exists only on the AP branch.

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 partitionAP — answer during the partition
What the minority side doesRejects writes, and usually reads tooServes and accepts writes locally
What the client seesErrors and timeouts, never a stale committed valueAn answer — possibly stale, possibly one of several versions
Work owed afterwardsNone; state never forkedReconciliation and compensation
Where it bitesAn invariant-free workload pays availability for nothingEvery invariant needs a merge rule, and some have none
When the cost endsOnce the link heals and a leader is re-elected — the election itself is the CP branch's lingering costLater: 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 R and W, 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.

Learning resources