CAP Theorem in Practice
What CAP really constrains (only behaviour during a partition), why "pick two" is wrong, and how PACELC, quorums and consistency models turn it into design decisions.
The CAP theorem says one narrow thing: when the network splits your replicas into groups that cannot talk, each request must either be answered from possibly stale local data (available) or refused until the split heals (consistent). It is not a menu where you pick two of three, and it says nothing about the 99.9% of the time when there is no partition. The useful design questions are what you do during a partition, what you trade for latency the rest of the time, and which consistency guarantee each feature actually needs.
Context
Eric Brewer stated CAP as a conjecture in a 2000 keynote; Seth Gilbert and Nancy Lynch proved a precise version in 2002, with "consistency" meaning linearizability and "availability" meaning every request to a non-failed node gets a non-error response. During the NoSQL wave of 2008-2013 the theorem became a marketing label: Amazon's Dynamo paper (2007) and its descendants Cassandra and Riak were sold as "AP", HBase and ZooKeeper as "CP", and relational databases were filed, wrongly, under "CA". Brewer himself wrote in 2012 that the "two out of three" framing was misleading; Daniel Abadi proposed PACELC the same year to add the latency trade-off that dominates normal operation; Martin Kleppmann's 2015 critique and Kyle Kingsbury's Jepsen tests (2013 onward) showed how far vendor labels were from measured behaviour.
You have already met the trade-off, probably by accident: a Postgres read replica that returned a row the user had just deleted, a DynamoDB call with the ConsistentRead flag, a Cassandra query with CONSISTENCY QUORUM, or etcd and Consul being used for leader election precisely because they refuse to answer when they cannot agree. Even S3 was eventually consistent for overwrites until December 2020.
The simplest example needs two nodes and a cut cable. A client writes x = 1 to node A. The link between A and B is down. Another client asks B for x. B can say 0 (available, wrong) or say "I cannot know right now" (consistent, unavailable). There is no third answer, and that is the whole theorem.
- Partition
- A network failure where some nodes cannot reach others but all keep running. Not a crash: both sides think they are fine.
- Linearizability
- The strong consistency CAP means: every read returns the latest completed write, as if there were one copy of the data.
- Availability (CAP sense)
- Every request to a live node eventually gets a real answer. Not "99.99% uptime", which is a different word for a different thing.
- Eventual consistency
- Replicas may disagree for a while; if writes stop, they converge. Says nothing about how long "a while" is or what a read sees meanwhile.
- Quorum
- The minimum number of replicas that must acknowledge a write or answer a read for it to count.
Why it matters
Every system with more than one copy of a piece of data has already made this choice, usually by default and usually without anyone noticing: read replicas, caches, multi-region deployments and CDNs all serve possibly stale data in exchange for answering fast and answering at all. The bugs are the "I just saved it and it is gone" class, two customers buying the last unit, and a counter that goes backwards, and they only appear under real traffic across real network hiccups.
What the theorem actually constrains
Three consequences follow from the precise definitions, and they are what separate a real answer from the slogan.
- You do not choose P. Partitions happen: a switch reboots, a cross-region link degrades, a garbage-collection pause makes a node look dead. A system that "chooses CA" is a single node, or a cluster that has simply not decided what it will do when the cable is cut and will do something inconsistent by accident.
- C and A are per operation, not per system. The same database can refuse writes during a partition and still serve reads, or be consistent for one table and available for another. Labels like "an AP database" describe a default, not a law.
- Nothing is said about the healthy case. With no partition, a system can be both consistent and available. What it cannot be is also fast, which is the trade-off the next section adds.
| Letter | Precise meaning | What people often think it means |
|---|---|---|
| C | Linearizability: reads reflect the latest completed write, system-wide | The C in ACID, or "no data corruption". Unrelated. |
| A | Every request to a non-failed node gets a non-error response, eventually | High uptime / HA (high availability). A CP system can have five nines; it just returns errors during partitions. |
| P | The system keeps operating when messages between nodes are lost | "Handles node crashes". Crashes are easy; a live node that cannot be reached is the hard case. |
PACELC: the trade-off you make every day
Abadi's extension reads: if there is a Partition, choose Availability or Consistency; Else, choose Latency or Consistency. The second half is the one that matters in production, because partitions are rare and latency is constant. To be consistent, a write must reach a majority (or all) of replicas before it is acknowledged, and a read must consult enough replicas to be sure it sees the latest write. Every such round trip is latency. To be fast, you acknowledge after one replica and read from the nearest one, and now replicas can disagree.
| System (typical default) | During partition | Normal operation | Why |
|---|---|---|---|
| Cassandra, DynamoDB (default eventual reads), Riak | A | L | Dynamo lineage: any replica can accept a write; conflicts resolved later. Built for always-on carts and feeds. |
| MongoDB (primary reads, majority writes) | A for reads on stale secondaries, C on primary | C | Single primary per shard; a partitioned primary steps down after election timeout. |
| PostgreSQL with synchronous replica | C (writes block) | C | A commit waits for the standby; if the standby is unreachable, commits hang or fail. |
| Spanner, CockroachDB, etcd, Consul, ZooKeeper | C | C | Consensus (Paxos/Raft): every write needs a majority. Minority side refuses. Latency is one quorum round trip. |
| Redis async replication, read replicas, CDN caches | A | L | Replicas lag by design; nobody waits for them. Stale reads are the accepted price. |
Quorums: consistency as a dial
Leaderless and multi-replica stores expose the dial directly. With N replicas, a write waits for W acks and a read consults R replicas. If W + R > N, every read set overlaps every write set in at least one replica, so the read sees the newest acknowledged value (the replica returns the version with the latest timestamp). Lower either number for speed and availability; the overlap disappears and stale reads become possible.
// DynamoDB: eventually consistent read (default, half the cost, may lag ~1s)
await ddb.send(new GetItemCommand({TableName: 'carts', Key: {id: {S: '42'}}}))
// Same read, linearizable: routes to the leader replica for that partition
await ddb.send(new GetItemCommand({
TableName: 'carts', Key: {id: {S: '42'}}, ConsistentRead: true,
}))
// Cassandra: per-statement consistency level
await client.execute(
'SELECT balance FROM accounts WHERE id = ?', [42],
{consistency: types.consistencies.quorum}, // majority of replicas
)
await client.execute(
'UPDATE accounts SET balance = ? WHERE id = ?', [90, 42],
{consistency: types.consistencies.quorum}, // W + R > N holds
)The ladder of consistency models
"Consistent" and "eventual" are the ends of a ladder, and the middle rungs are where most products actually live. Each rung is a promise a client can rely on; each costs something to provide.
| Model | Promise | Typical cost |
|---|---|---|
| Linearizable | One global order; a read after a write (by anyone) sees it | Quorum or leader round trip per operation; cross-region latency |
| Causal | If A happened before B, everyone sees A before B; unrelated writes may reorder | Version vectors or session tokens; cheap in one region |
| Read-your-writes | A client always sees its own previous writes | Route that client’s reads to the primary, or carry a write timestamp |
| Monotonic reads | A client never sees data go backwards in time | Sticky replica per session |
| Eventual | Replicas converge if writes stop | Almost nothing; conflict resolution is your problem |
Choosing per feature
The mistake is deciding once for the whole system. The right unit is the feature, and the question is "what is the cost of a stale answer, and what is the cost of no answer?".
| Feature | Stale answer costs | No answer costs | Pick |
|---|---|---|---|
| Shopping cart, likes, view counts | Almost nothing; merge later | Lost sales, angry users | Available. Accept writes anywhere; merge (union of cart items, sum of counts). |
| Last unit of inventory, seat booking | Two people get one seat | A user retries in a second | Consistent. Single writer per item, or a consensus store, or a DB row lock. |
| Account balance, payments | Money is created or lost | Payment delayed, retried idempotently | Consistent. Serializable transactions on one primary. |
| User profile after edit | User thinks the save failed | Rare | Eventual plus read-your-writes: read from primary for that user briefly after a write. |
| Leader election, locks, feature flags | Two leaders, double processing | Work pauses until quorum returns | Consistent. That is what etcd, Consul and ZooKeeper are for. |
Pitfalls
- Reading from a replica right after writing to the primary
The write is acknowledged, the next request is load-balanced to a replica that has not applied it yet, and the UI shows the old value or a 404. It is the most common consistency bug in ordinary web apps. Route reads to the primary for a short window after a write (a cookie with the write's log position or a timestamp), or return the written object in the response and render from it.
- Believing "eventual" converges on its own
Convergence needs a rule for concurrent writes. The default in most stores is last-writer-wins by timestamp, which silently discards one of two edits made in the same second on two replicas. If both edits matter, you need a merge (CRDTs, the union of a cart) or a version check that rejects the second write.
- Calling a single-primary database "CA"
A primary with async replicas serves stale reads during a partition (A) and cannot accept writes on the replica side (not A); a primary with sync replicas blocks commits when the replica is unreachable (not A). Either way P was never optional. Say what the system does during the partition instead of assigning letters.
- Even-sized quorum clusters
Four nodes need three for a majority, which is the same tolerance as three nodes (one failure) with more hardware and more ways to split evenly. Consensus clusters are sized 3, 5 or 7 for this reason.
- Confusing CAP availability with uptime
A consistent system that returns errors for the 30 seconds a partition lasts can still have 99.99% uptime; an available system can be down for a week for unrelated reasons. Uptime is an SLA (service level agreement) number about operations; CAP availability is a statement about behaviour when messages are lost. Mixing them produces designs that are slow all year to protect against a case that lasts seconds.
Interview questions
Q1Explain the CAP theorem without saying "pick two".
When replicas cannot communicate, a node that receives a request must either answer from what it has, which may be stale, or refuse until it can confirm with the others. That is the whole constraint, and it only applies while the partition lasts. Partitions are not something you opt out of, so the real choice is what each operation does during one: stay available with weaker consistency, or stay consistent and return errors on the minority side.
Q2Can a system be CA?
Only if it never has a partition, which means a single node. Any system with two nodes over a network will eventually lose messages between them, and at that moment it will behave as either AP or CP whether or not anyone designed it. A "CA" label usually means the behaviour under partition was never specified, which is the worst of the three options.
Q3What does PACELC add, and why does it matter more day to day?
It adds the trade-off with no partition: else, latency versus consistency. A consistent write needs a majority acknowledgement before returning and a consistent read needs to consult enough replicas, so every consistent operation pays a round trip, across regions if the replicas are spread out. Partitions last seconds a few times a year; that latency is paid on every request, so it is the trade-off that shapes the architecture.
Q4Design a shopping cart that works across two regions during a network split.
Make the cart available: each region accepts adds and removes locally and replicates asynchronously. Represent the cart so that concurrent changes merge, for example a set of items with add and remove timestamps, or an observed-remove set CRDT (conflict-free replicated data type), so a user who added on one side and removed on the other converges to a sensible cart. Move the consistency requirement to checkout, which runs against a single consistent inventory and payment store and can fail cleanly if that store is unreachable.
Q5What happens during a partition in Postgres with a synchronous replica, versus in Cassandra?
Postgres with synchronous replication will not acknowledge a commit until the standby confirms it, so if the standby is unreachable, writes hang until timeout or an operator switches the mode; reads on the primary keep working. That is choosing consistency. Cassandra with default consistency ONE keeps accepting writes and reads on both sides of the split, storing hints to replay later, and the two sides can hold different values until repair. That is choosing availability, and a QUORUM setting would make the minority side refuse instead.
Q6A user saves their profile and the next page shows the old data. What is going on and how do you fix it?
The write went to the primary and the next read was served by a replica that had not applied it yet, so the system is eventually consistent but the feature needs read-your-writes. Fixes, cheapest first: render the saved object from the write response instead of re-fetching; route that user's reads to the primary for a few seconds after a write, tracked by a cookie or session flag; or have the client send the write's log position and let the replica wait until it has caught up.
Q7Why does W + R > N give you the latest write, and what does it not give you?
Any read set of R replicas must share at least one member with any write set of W replicas when W plus R exceeds N, so at least one replica the read consults has the newest acknowledged value, and the client picks the newest version among the answers. It does not give full linearizability by itself: two concurrent writes can be partially applied, a failed write can be visible on one replica, and clock-based versioning can pick the wrong winner. It rules out the simple stale read, not every anomaly.
- CAP constrains one thing: during a partition, each operation is either available (may be stale) or consistent (may refuse). Partitions are not optional, so "CA" is not a choice.
- C means linearizability and A means every live node answers. Neither is ACID consistency or uptime.
- PACELC adds the trade-off that runs all year: consistency costs a quorum round trip on every operation, availability and low latency cost stale reads.
- Consistency is a dial per operation (quorums, read concerns, consistency levels), and a ladder of models: linearizable, causal, read-your-writes, monotonic, eventual.
- Decide per feature by the cost of a stale answer versus no answer: carts and counters go available, money and uniqueness go consistent, profiles go eventual plus read-your-writes.
- "Eventual" needs a merge rule; last-writer-wins silently drops data. Size consensus clusters odd. Trust Jepsen over the marketing page.
- Gilbert & Lynch (2002) — Brewer’s conjecture and the feasibility of consistent, available, partition-tolerant web services
- Eric Brewer (2012) — CAP Twelve Years Later: How the "Rules" Have Changed
- Daniel Abadi (2012) — Consistency Tradeoffs in Modern Distributed Database System Design (PACELC)
- Martin Kleppmann (2015) — A Critique of the CAP Theorem
- Jepsen — distributed systems safety analyses