The CAP theorem says that when the network between the machines in a distributed system breaks, each machine has to choose: answer with data that might be out of date, or refuse to answer until it can be sure. Knowing which one your system picks tells you exactly how it will misbehave on a bad day, which is why it turns up in so many system design interviews.
The three letters
A distributed system is several machines, called nodes, that hold copies of the same data and talk to each other over a network. CAP names three properties you would like it to have.
Consistency
Every read sees the most recent write, as if there were only one copy of the data. If you update your address and then read it back from any node, you get the new address. The formal name for this is linearisability.
This is not the C in ACID. ACID's consistency means a transaction leaves the database obeying its rules (no negative balances, no orphaned rows). CAP's consistency is about copies agreeing with each other.
Availability
Every request that reaches a working node gets a real answer, not an error or a timeout. There is no promise that the answer is the latest one, only that there is one.
This is also a narrower idea than everyday 'availability'. A system can have excellent uptime and still not be available in the CAP sense, because it deliberately turns some requests away during a network fault.
Partition tolerance
A partition is a network fault that splits the nodes into groups that can't reach each other. Messages between the groups are lost or delayed indefinitely. Partition tolerance means the system keeps running, in some form, while that is going on.
Why partition tolerance isn't optional
Networks fail. A switch dies, a cable is cut, a cloud region loses its link to another, or a node pauses for so long that its peers give up on it. From the inside, a node can't tell a slow peer from a dead one or from a broken link: all it sees is silence.
So if you run more than one machine and they talk over a network, partitions will happen whether you plan for them or not. A single database server never has a partition, but it isn't distributed either, and it has no copy to fall back on when it fails. That is why 'pick any two of the three' is misleading. For a real distributed system, P comes with the territory, and the choice that matters is between C and A while a partition lasts.
The real choice: consistency or availability
Picture two nodes, A and B, each holding a copy of a shop's stock count. The link between them breaks. A customer buys the last item through node A, so A now says zero. Then another customer asks node B how many are left.
Node B has not heard about the sale, and it can't: the message has nowhere to go. It has exactly two options. It can answer with what it has, which is now wrong, or it can refuse until it can check with A. There is no third option, because no amount of cleverness lets B know about a write it never received.
Here is that sequence, and the two ways B can respond:
One read during a network partition
Step 1 of 7: Both nodes hold a copy of the stock count: one item left.
In code, the decision sits in each node's read path. A rough sketch:
def read(key):
if can_reach_majority():
return latest_value(key)
if MODE == "CP":
# refuse rather than risk stale data
raise Unavailable("try again later")
# AP: answer now, maybe out of date
return local_value(key)Choosing consistency (CP)
A CP system keeps its answers correct by making some requests wait or fail. Usually the group of nodes that still holds a majority carries on, and the smaller group stops accepting writes, and often reads too, until the partition heals. Systems built on consensus, such as etcd and ZooKeeper, work this way: nothing is accepted without a majority agreeing.
The cost is visible to users: errors, timeouts and 'please try again' for anyone whose request lands on the wrong side. It fits data where a wrong answer is worse than no answer, such as account balances, the last item in stock, locks and deciding which node is the leader.
Choosing availability (AP)
An AP system keeps every node answering and accepting writes, then reconciles once the nodes can talk again. Dynamo-style stores such as Cassandra lean this way by default.
The cost is less visible but real. Readers can see stale data, and two sides can accept conflicting writes to the same record. Something has to settle those conflicts later: 'last write wins' is simple but silently throws one of the writes away, while merging (say, combining two versions of a shopping basket) keeps both but needs design work. It fits data where a slightly old answer is harmless, such as like counts, feeds, reviews and recommendations.
A worked example: one shop, two choices
An online shop runs in two data centres, London and Frankfurt, each with a full copy of its data. One afternoon the link between them drops for ten minutes, and customers keep arriving on both sides.
The shop does not have to make one choice for everything. It can decide per kind of data:
| Data | Choice | Why |
|---|---|---|
| Stock of rare items | CP | Selling one item twice means refunds |
| Payments | CP | A wrong balance is worse than a delay |
| Reviews | AP | Ten-minute-old reviews hurt no one |
| Basket | AP | Merge both sides when the link returns |
So a Frankfurt customer might see 'checkout temporarily unavailable' for a rare item, while browsing, reading reviews and filling a basket all carry on as normal. Most real systems look like this: a mix, chosen by what a wrong answer would cost.
When there is no partition
Most of the time the network is fine, and CAP says nothing about that case. There is still a trade-off, though. Keeping copies in step means waiting for other nodes to confirm each write, and waiting takes time.
The PACELC model extends CAP to cover this: if there is a Partition, choose Availability or Consistency; Else, choose Latency or Consistency. Many databases let you make that choice per request. Cassandra, for example, lets a read ask one replica (fast, possibly stale) or a majority of them (slower, up to date).
Common mistakes
- Reading it as 'pick any two'. A distributed system can't opt out of partitions, so the real decision is what to give up while one lasts.
- Mixing up the two Cs. CAP's consistency is about copies agreeing; ACID's is about a transaction keeping the data valid.
- Treating availability as uptime. CAP availability is a strict promise about every request during a fault, not a yearly percentage.
- Labelling a whole database CP or AP. Many databases are configurable, and different operations in one system can behave differently.
- Forgetting timeouts. A node only knows about a partition because a peer went quiet for too long. Set that limit too short and you trigger the trade-off on a merely slow network.
Answering it in an interview
Define the three properties precisely, including that CAP's consistency is not ACID's. Say that partitions are a fact of life, so the question is consistency or availability during one. Then make the choice for the specific data in the question, name what it costs the user, and mention PACELC to show you know the trade-off does not stop when the network is healthy.
Key takeaways
- CAP is about what a distributed system does during a network partition.
- Partitions will happen, so the real choice is between consistency and availability.
- CP returns errors or waits rather than serve stale data; AP always answers but may be out of date and must reconcile later.
- Choose per kind of data: money and scarce stock lean CP, feeds and counters lean AP.
- With no partition, the trade-off becomes latency versus consistency (PACELC).