Understanding the CAP Theorem in Distributed Systems
May 1, 2023 · 2 min read · Distributed Systems · CAP Theorem · Architecture
Eric Brewer introduced the CAP theorem in 2000. It says a distributed data system can guarantee only two of three properties at the same time: consistency, availability, and partition tolerance.
- Consistency: every read returns the most recent write, or an error.
- Availability: every request receives a response, even if the data is not the most recent.
- Partition tolerance: the system keeps operating while network communication between nodes is broken.
Why you can only pick two
During a partition, nodes cannot reach each other. A node that receives a request has two options: answer with the data it has, or refuse to answer until it can confirm the data is current. It cannot do both. Partitions are not a question of if but when, so practical systems assume partition tolerance and choose between consistency and availability.
CP and AP in practice
A CP system may reject requests during a partition rather than serve stale data. An AP system keeps responding, possibly with data that is out of date. The databases most teams use sit on one side of this choice:
- CP: HBase, BigTable, Zookeeper, etcd
- AP: MongoDB, CouchDB, DynamoDB (with tunable consistency)
Choosing per operation
The right trade-off depends on the operation. A bank transfer should be CP: better to return an error than to let the same balance be spent twice. An e-commerce cart can be AP: a user adds an item and it succeeds even when the backend cannot confirm the latest state right away.
CAP is a model, not a law. Systems with tunable consistency pick a level per operation — strong consistency for payments, eventual consistency for analytics — so the choice is rarely a single setting for the whole system.