The CAP Theorem
In 2000, computer scientist Eric Brewer conjectured — and two years later formally proved — that a distributed data store can satisfy at most two of three properties simultaneously:
- Consistency
- Availability
- Partition Tolerance
This is the CAP Theorem, one of the most cited — and most misunderstood — ideas in distributed systems.
Defining the Three Properties
Consistency
Every read receives the most recent write or an error. All nodes in the system see the same data at the same time. If you write x = 5 on Node A, a subsequent read of x on Node B must return 5.
Note: this is not the same as the "C" in ACID. ACID consistency means the database moves from one valid state to another. CAP consistency means all nodes agree on the current value.
Availability
Every request receives a response — not necessarily the most recent data, but a response (not an error, not a timeout). An available system keeps responding even when some nodes are down.
Partition Tolerance
The system continues to operate despite network partitions — situations where messages between some nodes are lost or delayed indefinitely. A partition could be two data centres losing connectivity, or a network switch failing in a cluster.
Why You Can Only Pick Two
Consider a simple scenario: two nodes (A and B) and a client.
Scenario: A network partition occurs — Node A and Node B cannot communicate.
- A client writes
x = 5to Node A. - A second client reads
xfrom Node B.
Now the system must choose:
- Prioritise Consistency: Node B refuses to respond until it can synchronise with Node A (which may be impossible during the partition). The system is consistent but unavailable.
- Prioritise Availability: Node B returns its stale value of
x(maybex = 3from before the partition). The system is available but inconsistent.
There is no third option. This is the fundamental tension.
Because network partitions are not optional (real networks fail), the practical choice is between CP and AP:
- CP systems sacrifice availability during partitions (they error or block until the partition heals).
- AP systems sacrifice consistency during partitions (they return potentially stale data).