← All posts
ConceptsJuly 22, 202610 min read

The CAP Theorem Explained (for System Design Interviews)

The CAP theorem is the most quoted and most misquoted result in distributed systems. "Pick two of three" is the popular summary, it is on every interview cheat sheet, and it is wrong in a way that actively misleads people designing systems.

This post covers what CAP actually claims, why the choice is narrower than three-choose-two, what a client genuinely observes on each side of it, why PACELC matters more day to day, and the nuance that separates a memorised answer from an understood one — that the choice is made per operation, not per database.

The three words, precisely

CAP describes a system storing data across multiple nodes.

Consistency (C) means linearizability: every read returns the most recent completed write, as though there were a single copy of the data. Note this is a much stronger claim than the "C" in ACID, which only means the database's own constraints and invariants hold. Same letter, different property — and mixing them up is a common interview stumble.

Availability (A) means every request sent to a non-failing node receives a non-error response. This definition is stricter than the colloquial sense of "available". It is not "99.9% uptime". It is: no node that is still running is allowed to return an error or hang, ever, even when it has no idea what the rest of the cluster is doing.

Partition tolerance (P) means the system keeps operating when the network drops or arbitrarily delays messages between nodes.

Being precise about A is worth doing out loud, because under CAP's definition almost no real system is strictly "available" — most will happily return a 503 somewhere. The theorem is a formal result about a formal model, not a product description.

Why "pick two of three" is wrong

The phrase implies you sit down and choose any two. You don't, because P is not a choice you get to make.

majority side n1 n2 n3 minority side n4 n5 network partition every node is healthy — only the link is gone client client
A partition is not a node failure. Every machine here is up and serving. The network between them simply stopped delivering messages, and neither side can tell whether the other is dead or merely unreachable. That ambiguity is the entire problem.

Networks partition. Switches fail, cables get unplugged, a cloud availability zone loses connectivity, a GC pause makes a node look gone. You cannot opt out of that any more than you can opt out of gravity. So "CA" — a system that gives up partition tolerance — is not a design choice; it is a system that produces wrong answers when the network misbehaves, which it eventually will.

The real question is narrower and sharper: when a partition happens, and a node cannot reach its peers, does it answer anyway or does it refuse?

What each choice actually looks like to a client

This is where the abstraction becomes concrete, and where interviewers probe.

CP — choose consistency, sacrifice availability majority (n1 n2 n3) has quorum · accepts writes minority (n4 n5) no quorum · refuses write → 503, try later. The data is never wrong; it is briefly unreachable. AP — choose availability, sacrifice consistency majority (n1 n2 n3) x = 7 minority (n4 n5) x = 9 write → 200 OK. Two truths now exist, to be reconciled on heal.
Same failure, two behaviours. CP turns a network problem into an error your application must handle. AP turns it into a correctness problem your application must reconcile. Neither makes the partition go away — they choose which kind of pain to absorb, and where.

The AP path has a consequence people skip: once two sides accept conflicting writes, something has to resolve them. Last-write-wins is the usual default and it silently discards data — two users edit the same field, one edit evaporates. The alternatives are vector clocks with application-level merge logic (Dynamo's original approach, which pushed the problem to the caller), or CRDTs, which are data structures designed so that concurrent updates merge deterministically. Saying "AP, and I'd resolve conflicts with X" is a much stronger answer than "AP, eventual consistency" full stop.

The 99.9% of the time nothing is partitioned

Here is the part the popular summary omits entirely: when the network is healthy, CAP forces no trade-off at all. You get consistency and availability together. Partitions are rare — a well-run cluster might see minutes of them a year.

So a theorem about the rare case tells you almost nothing about the common case. What governs the common case is latency, and that is what PACELC adds:

is there a partition? yes no rare — minutes per year Availability or Consistency the other 99.9% of the time Latency or Consistency CAP only ever spoke about this branch. Strong consistency costs a round trip to other replicas on every operation. That bill arrives constantly, not rarely.
PACELC: if Partition, then A or C; Else, L or C. The "else" branch is the one your users actually feel. Keeping replicas linearizable means coordinating across nodes before answering, and that coordination is measured in milliseconds on every single request — which is why many systems relax consistency in normal operation, not just during failures.

Bringing up PACELC unprompted signals that you understand consistency as a spectrum with a continuous latency price, rather than a switch that only flips during outages.

What a partition actually looks like from the inside

Theory gets slippery here, so it helps to walk a concrete timeline. A five-node cluster, replication factor three, and a switch fails between two racks.

t+0s. The link drops. Nothing detects anything yet. Both sides continue serving as though the cluster is whole, because neither has any reason to suspect otherwise. This window is short but it is real, and writes accepted during it are the ones most likely to conflict.

t+3s. Heartbeats start timing out. Each side now observes that two peers have gone quiet — and crucially, neither side can tell whether those peers crashed or are merely unreachable. From n1's point of view, "n4 is dead" and "I am cut off from n4" produce identical evidence.

t+5s. The failure detector fires and the system has to act on a guess. In a quorum-based design, the majority side notices it still has three of five and elects to continue; the minority side notices it has only two and stops accepting writes. In a leaderless AP design, both sides simply keep going.

t+5s to t+40s. This is the window everything is decided in. A CP system is now returning errors to every client that happens to be routed to the minority side — users see failures even though the machines serving them are perfectly healthy. An AP system is happily accepting writes on both sides, and divergence accumulates in proportion to traffic.

t+40s. The switch is replaced and the partition heals. A CP system re-syncs the minority from the majority's log and resumes; there is nothing to reconcile because the minority never accepted anything. An AP system must now merge two histories, and every conflicting key needs a resolution rule — which is precisely the work the AP choice deferred rather than avoided.

Three things are worth taking from that timeline. The choice is made by a timeout, not by knowledge. The cost of CP is paid by users during the partition; the cost of AP is paid by engineers after it. And "eventual consistency" has a duration attached — eventual means "once the partition heals and the merge completes", which could be forty seconds or could be the next morning if nobody noticed.

Split brain, and why fencing matters

The most dangerous outcome is not either choice made cleanly. It is split brain: both sides believing they are authoritative at once.

Imagine the partition above, but the minority side wrongly concludes the majority is dead and elects its own leader. Now there are two leaders, both accepting writes, both convinced they are the only one. Every guarantee the system offered is void.

Quorums prevent the common case — a group of two out of five cannot form a majority, so it cannot elect a leader. But quorums do not cover the nastiest variant, which is a leader that was legitimately elected, then paused (a long GC, a VM suspend), had its lease expire, watched a new leader get elected, and then woke up still believing it holds the lock. It is not malicious and it is not partitioned any more. It is simply stale, and it is about to write.

The defence is a fencing token: a monotonically increasing number handed out with the lock, which the storage layer checks and rejects if it is lower than the highest it has already seen. The woken-up old leader presents token 33, the storage sees it already accepted token 34, and the write is refused. Without fencing, a lock service gives you a guarantee that a sufficiently long pause can silently break — which is a good thing to be able to say when an interviewer asks "what if the leader freezes?"

It is a per-operation choice, not a per-database one

This is the nuance that most reliably separates candidates, because the cheat-sheet version implies you classify a database as "CP" or "AP" once and move on. Real systems are tunable, usually per query.

Dynamo-style stores expose N (replicas), W (replicas that must acknowledge a write) and R (replicas that must respond to a read). The rule is straightforward arithmetic:

N = 3   replicas of each key

W = 1, R = 1   →  R + W = 2 ≤ N   fast, may read stale data      (AP-ish)
W = 2, R = 2   →  R + W = 4 > N   read and write sets overlap    (CP-ish)
W = 3, R = 1   →  R + W = 4 > N   slow writes, very fast reads

When R + W > N, the set of replicas you read from must overlap the set that acknowledged the write, so you are guaranteed to see it. That is one line of arithmetic, and it is worth being able to derive on a whiteboard rather than recite.

Cassandra exposes exactly this as per-query consistency levels (ONE, QUORUM, ALL). MongoDB has write concerns and read concerns. DynamoDB lets you ask for an eventually consistent or a strongly consistent read on the same table, at different prices.

So the strong answer is per-workload: "The session store runs at ONE — a stale read is harmless and latency matters. The billing ledger runs at QUORUM — I'd rather fail the write than double-charge." Same cluster, two different points on the curve.

Where systems do sit firmly: coordination services like etcd and ZooKeeper are unambiguously CP, because their entire purpose is to be the thing everyone agrees with — an available-but-wrong lock service is worse than useless. DNS is famously AP and nobody minds, because stale records resolving for a few minutes is an acceptable cost for never being down.

What CAP does not tell you

Naming the limits is the fastest way to show you have thought past the cheat sheet.

It says nothing about how long a partition lasts. A system that refuses writes for 200 milliseconds and one that refuses them for six hours are both "CP". The operational difference is enormous and CAP is silent on it.

It says nothing about durability. A system can be perfectly consistent and available and still lose your data on a power cut. That is the D in ACID, and it is orthogonal.

Its "availability" is not your SLA. CAP's A is an absolute, formal property. Your uptime target is a statistical one. A system can violate CAP-availability constantly and still hit four nines.

It assumes an asynchronous network with no clocks. Real networks have bounded-ish latency and roughly synchronised clocks, which is precisely what systems like Spanner exploit — using tightly synchronised time to offer external consistency and very high availability at once. That does not break the theorem; it sidesteps the model's assumptions with hardware.

Partition detection is itself a guess. A node cannot distinguish "peer is dead" from "peer is slow". Every CP system is really making a timeout-based bet, and tuning that timeout trades false failovers against slow ones.

Using it in an interview

Don't recite the definition. Apply it per datastore, in one sentence each, as you introduce them:

  • "The payments ledger needs linearizable writes — CP. During a partition I'd rather reject the transaction than risk a double charge, and I'd surface that to the user as a retry."
  • "The activity feed is AP. A few seconds of staleness is invisible; being down is not. I'd resolve concurrent writes last-write-wins, since a feed entry is immutable once written."
  • "Both live in the same Cassandra cluster, at different consistency levels."

The traps to avoid: saying "CA" as though it were an option; using CAP's "availability" as a synonym for uptime; classifying a whole database as CP or AP when it is tunable; and — most common — describing the theorem accurately and then never applying it to the system on the whiteboard. The definition earns no points. The application does.

It pairs naturally with two other decisions you will be asked about in the same breath: which store to reach for in the first place, covered in SQL vs NoSQL, and how data is spread across nodes so that a partition only takes out part of it, covered in consistent hashing.

If you want to find out whether this holds up under follow-up questions — "what does the client see, exactly?" is the usual one — you can run a full system design mock by voice on Whitepad.

Practice this out loud

Reading is the easy part. Sit across from a senior AI interviewer that talks, watches your whiteboard, and scores you like the real thing — your first mock is free.

Start a free mock →