CAP Theorem & Consistency Models
Once data lives on more than one machine, you face trade-offs that don't exist on a single server. The CAP theorem captures the sharpest one: when the network partitions (nodes can't talk to each other), a distributed system must choose between consistency (every read sees the latest write) and availability (every request gets a response). PACELC extends it: even without partitions, you trade latency against consistency.
These ideas are central to system design interviews and, more importantly, to choosing databases and designing replication. The practical skill isn't reciting "pick two of three". It's knowing which consistency model each part of your system actually needs, and what the database you pick really guarantees.
TL;DR
- CAP: during a network partition, a system can't be both linearizably consistent and fully available. It must refuse some requests (CP) or serve possibly stale data (AP).
- Partitions aren't optional in distributed systems, so the real choice is C vs A when partitions happen.
- PACELC: if Partitioned, choose A or C; Else, choose Latency or Consistency.
- Consistency is a spectrum: linearizable → sequential → causal → read-your-writes and monotonic reads → eventual.
- Quorums (R + W > N) give tunable consistency in replicated stores.
- Choose per use case: payments and inventory need strong guarantees; feeds, counters, and caches tolerate eventual consistency.
Quick Example
A replicated key-value store with N = 3 replicas, tuning consistency per request with quorums:
During a partition that isolates one replica, QUORUM operations still succeed on the majority side, while the minority side can't satisfy QUORUM: availability is sacrificed there to keep consistency.
Core Concepts
The CAP Theorem
Formalized by Gilbert and Lynch (2002) from Eric Brewer's conjecture:
- Consistency (C): specifically linearizability. Every read returns the most recent write, as if there were a single copy.
- Availability (A): every request to a non-failing node receives a (non-error) response.
- Partition tolerance (P): the system keeps operating despite dropped or delayed messages between nodes.
When a partition splits nodes, a node receiving a write can either reject it or wait until it can coordinate (keeping C, losing A), or accept it and let replicas diverge (keeping A, losing C). Since networks will partition, "CA" isn't a real option for distributed systems. You're choosing behavior during partitions.
Common Misreadings
- CAP says nothing about behavior when there's no partition, which is most of the time. That's PACELC's territory.
- CAP's "C" is linearizability, not the "C" in ACID (database transactions).
- Systems aren't simply "CP" or "AP". Many are configurable per operation (quorum levels), and many provide neither full C nor full A under failure.
PACELC
If Partition: Availability or Consistency; Else: Latency or Consistency. Strong consistency requires coordination (waiting for replicas or consensus), which adds latency on every request, even when the network is healthy. Examples:
Consistency Models
From strongest to weakest:
Session guarantees (read-your-writes, monotonic reads) often give users a correct-feeling experience at far lower cost than global linearizability. See database replication.
Quorums
With N replicas, requiring W acknowledgements for writes and R responses for reads gives overlapping read and write sets when R + W > N, so reads see the latest acknowledged write (subject to details like sloppy quorums and clock skew). Consensus protocols such as Raft and Paxos build linearizable systems on majority quorums. See distributed consensus.
Handling Divergence in AP Systems
Available-under-partition systems accept conflicting writes and must reconcile them: last-write-wins by timestamp (simple, but loses updates), vector clocks to detect concurrent versions, application-level merges, or CRDTs (conflict-free replicated data types) that merge automatically, as used in collaborative and local-first apps.
Best Practices
Decide Consistency per Operation
Not all data needs the same guarantees. Use strong consistency for money, inventory reservations, and uniqueness constraints, and eventual consistency for feeds, recommendations, view counts, and caches.
Know Your Database's Real Guarantees
Read the documentation and Jepsen analyses: default isolation levels, replica read behavior, and failover data-loss windows. "Strongly consistent" marketing claims often come with configuration caveats.
Design the User Experience Around Staleness
Where eventual consistency is acceptable, show it honestly (for example "updated a few seconds ago"), route a user's reads to the primary right after their own writes, and avoid UIs that make stale data look authoritative.
Prefer Idempotent, Commutative Operations
Operations that can be safely retried and reordered (set membership adds, increments via CRDT counters, idempotent requests) make eventual consistency much easier to live with.
Common Mistakes
"We Chose CA"
A distributed database can't avoid partitions, so claiming CA just means behavior under partition hasn't been considered. Single-node databases are CA only because they aren't distributed.
Reading From Replicas After Writes
Users see their change "revert". Read from the primary for a short window after writes, or use session consistency tokens.
Using Last-Write-Wins for Important Data
Clock skew between nodes can make an older write "win", silently discarding newer data. Use version checks, consensus-backed stores, or CRDTs for data where lost updates matter.
FAQ
What does CAP theorem actually say?
In a distributed system experiencing a network partition, you can't guarantee both linearizable consistency and availability for every request. You must either reject or delay some requests to stay consistent, or serve them and risk returning stale or conflicting data.
Is MongoDB CP or AP? What about Cassandra?
It depends on configuration. MongoDB with majority write and read concerns behaves as CP: minority partitions can't accept writes. Cassandra is often described as AP, but with QUORUM reads and writes it provides much stronger consistency, trading availability during partitions. Most modern databases let you tune the trade-off.
What's the difference between strong and eventual consistency?
Strong (linearizable) consistency means every read reflects the latest completed write, as if there were one copy. Eventual consistency means replicas may temporarily disagree, but converge when updates stop. Strong consistency costs latency and availability; eventual consistency is faster and more available, but requires handling stale reads and conflicts.
How is PACELC different from CAP?
CAP only describes trade-offs during network partitions. PACELC adds that in normal operation you still trade latency for consistency, because coordinating replicas takes time. It better explains everyday database behavior and design choices.
Related Topics
- System Design — Designing large-scale systems
- Distributed Consensus — How strong consistency is achieved
- Database Replication — Synchronous vs asynchronous copies
- Consistent Hashing — Partitioning data across nodes
- High Availability — Staying up through failures
- Local-First — CRDTs and offline-first consistency