The CAP theorem and PACELC
What CAP actually says, why "pick two of three" is misleading, PACELC and the latency trade-off you face every day, where real databases sit, and how to use it in an interview without hand-waving.
Reading is half of it. See this used in a real interview: walk through Design a Distributed Key-Value Store →
CAP is the most quoted and most misquoted idea in distributed systems. Candidates often say "we choose AP" as if it were a setting. Interviewers want to hear that you know what the theorem actually constrains, that it only bites during a network partition, and that the trade-off you face every other day is between consistency and latency.
What CAP says
In a distributed data store, when a network partition separates some replicas from others, each side must choose between:
- Consistency (in CAP’s sense, linearizability): every read sees the most recent write, as if there were one copy.
- Availability: every request to a non-failed node gets a non-error response.
You cannot have both during the partition. If a replica cannot reach the others, it either refuses (stays consistent, loses availability) or answers with what it has (stays available, may be stale or accept conflicting writes).
That is all it says. Three clarifications matter:
- Partitions are not optional. Networks fail, so "CA" is not a real choice for a distributed system; you decide what happens when, not if, a partition occurs.
- It is per operation, not per database. One system can make balance updates consistent and profile reads available.
- The C is narrow. CAP’s consistency is linearizability, not ACID’s consistency. Plenty of useful guarantees (read-your-writes, causal consistency) sit in between.
PACELC: the everyday trade-off
Partitions are rare. Most of the time the network is fine, and the trade-off you actually make is latency. PACELC extends CAP:
If there is a Partition, choose Availability or Consistency; Else (normal operation), choose Latency or Consistency.
Keeping replicas strongly consistent means waiting for other replicas (often in other zones or regions) on every write, and sometimes every read. Relaxing consistency lets a node answer from its local copy. A cross-region consistent write costs a round trip of 50 to 150 ms; a local one costs 1 ms.
| System | During a partition | Normally | Shorthand |
|---|---|---|---|
| Single-leader Postgres or MySQL with async replicas | minority side cannot write | reads from replicas may be stale | PC/EC for the leader, EL if you read replicas |
| Spanner, CockroachDB | minority side unavailable | pay consensus latency on writes | PC/EC |
| Cassandra, DynamoDB (default) | both sides keep serving | low-latency local reads and writes | PA/EL, tunable per request |
| DynamoDB strongly consistent reads | read from the leader replica | higher read latency and cost | per-request EC |
| Redis Cluster | minority masters stop taking writes after a timeout, yet acknowledged writes can still be lost on failover | fast, asynchronous replication | neither cleanly: treat it as a cache, not a source of truth |
| etcd, ZooKeeper | minority unavailable | consensus on writes | PC/EC |
Choosing per piece of data
The practical way to use CAP in a design is to go through the data and ask what a stale read or a lost write costs.
- Money, inventory at checkout, unique usernames, seat holds: a wrong answer causes a wrong action that is hard to undo. Choose consistency; accept that writes fail during partitions and cost more latency.
- Feeds, likes, view counts, profiles, presence: a few seconds of staleness is invisible. Choose availability and low latency.
- Messages and shopping carts: must not be lost but can be merged; choose availability with conflict resolution (version vectors, CRDTs). See replication and consistency.
Saying "the order total and the payment state use a consistent store; the product catalogue and recommendations can serve stale reads" is a far better answer than picking a letter for the whole system.
Consistency is a spectrum
Between linearizable and "eventually, maybe" sit models that are often exactly what users need:
- Read-your-writes: a user always sees their own changes (route their reads to the leader for a short window, or carry a version token).
- Monotonic reads: a user never sees time go backwards (stick them to one replica).
- Causal consistency: if B replied to A, nobody sees B without A.
- Bounded staleness: reads are at most N seconds behind.
Most products need read-your-writes and causal ordering in a few places, not global linearizability.
What a partition looks like in practice
- Two regions lose their link: each sees the other as down.
- A zone’s switch fails, isolating a third of the nodes.
- A long garbage-collection pause makes a leader look dead; a new leader is elected; the old one wakes up still believing it leads. Fencing tokens and leases exist for exactly this; see consensus, leases and coordination.
The design question is always: on the minority side, do we refuse writes (consistent) or accept them and reconcile later (available)?
In the interview
- Do not lead with "this is an AP system". Lead with the data and what staleness costs.
- Say what happens during a partition for the important operations, and what the normal-case latency is.
- Name the mechanism: a single leader or consensus for the consistent parts; quorums, version vectors or CRDTs for the available ones.
Checklist
- Which data must be linearizable, and why.
- What each important operation does during a partition.
- The latency you pay for consistency in normal operation (PACELC).
- The weaker guarantees you do provide (read-your-writes, causal).
- How conflicts are resolved where you chose availability.