Topics
System Design

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.

Intermediate·14 min read·Updated Sep 28, 2026

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

client 1client 2write x = 1read x ?node A · x = 1node B · x = 0✕ partitionreplication cannot get throughanswer "x = 0"available · stalerefuse (503 / retry)consistent · unavailableA has the same problem in reverse: it cannot know whether B accepted a conflicting write.
The only situation CAP is about. B has not heard about the write. Answering is available but stale; refusing is consistent but unavailable. Nothing lets B answer correctly.

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.
LetterPrecise meaningWhat people often think it means
CLinearizability: reads reflect the latest completed write, system-wideThe C in ACID, or "no data corruption". Unrelated.
AEvery request to a non-failed node gets a non-error response, eventuallyHigh uptime / HA (high availability). A CP system can have five nines; it just returns errors during partitions.
PThe 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 partitionNormal operationWhy
Cassandra, DynamoDB
(default eventual reads), Riak
ALDynamo 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 primaryCSingle primary per shard; a partitioned primary steps down after election timeout.
PostgreSQL
with synchronous replica
C (writes block)CA commit waits for the standby; if the standby is unreachable, commits hang or fail.
Spanner, CockroachDB,
etcd, Consul, ZooKeeper
CCConsensus (Paxos/Raft): every write needs a majority. Minority side refuses. Latency is one quorum round trip.
Redis async replication,
read replicas, CDN caches
ALReplicas 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.

write, W = 2
r1 ✓r2 ✓r3 (async)ack after two
read, R = 2
r1r2 ✓ newestr3 ✓ staleoverlap on r2 → returns newest
W = 1, R = 1
write → r1read ← r3no overlap → stale is possible
N = 3. With W = 2 and R = 2 at least one replica (here r2) is in both sets, so the read cannot miss the write. With W = 1 and R = 1 the sets can be disjoint.
consistency.ts (DynamoDB and Cassandra, same dial)
// 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.

ModelPromiseTypical cost
LinearizableOne global order; a read after a write (by anyone) sees itQuorum or leader round trip per operation; cross-region latency
CausalIf A happened before B, everyone sees A before B; unrelated writes may reorderVersion vectors or session tokens; cheap in one region
Read-your-writesA client always sees its own previous writesRoute that client’s reads to the primary, or carry a write timestamp
Monotonic readsA client never sees data go backwards in timeSticky replica per session
EventualReplicas converge if writes stopAlmost 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?".

FeatureStale answer costsNo answer costsPick
Shopping cart, likes, view countsAlmost nothing; merge laterLost sales, angry usersAvailable. Accept writes anywhere; merge (union of cart items, sum of counts).
Last unit of inventory,
seat booking
Two people get one seatA user retries in a secondConsistent. Single writer per item, or a consensus store, or a DB row lock.
Account balance, paymentsMoney is created or lostPayment delayed, retried idempotentlyConsistent. Serializable transactions on one primary.
User profile after editUser thinks the save failedRareEventual plus read-your-writes: read from primary for that user briefly after a write.
Leader election,
locks, feature flags
Two leaders, double processingWork pauses until quorum returnsConsistent. 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.

Key takeaways
  • 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.

Preparing for interviews? DevRecall turns a job description into a prep plan that points at topics like this one.

Start free