The short answer
Quick answer: The CAP theorem says a distributed data store cannot guarantee all three of consistency (every read sees the latest write), availability (every request to a working node gets a response) and partition tolerance (the system keeps operating when the network between nodes fails). Since network failures will happen, partition tolerance is not optional. So the real choice is this: during a partition, do you refuse some requests to stay consistent (CP), or answer every request and risk returning stale or conflicting data (AP)?
The three properties, precisely
The everyday meanings of these words are looser than the theorem's.
- Consistency (C). This means linearizability: the system behaves as if there were a single copy of the data. Once a write completes, every later read, from any node, returns that value or a newer one. It is not the same as the "C" in ACID.
- Availability (A). Every request received by a node that has not crashed gets a non-error response. Not "most requests", and not "eventually": every one.
- Partition tolerance (P). The system continues to work even if messages between nodes are lost or delayed indefinitely.
The idea was proposed by Eric Brewer in 2000 and proved formally by Seth Gilbert and Nancy Lynch in 2002. The Wikipedia article has the history.
Why you cannot have all three
Take two nodes, A and B, each holding a copy of the data. The network link between them breaks. That is a partition.
- A client writes
x = 5to node A. - Node A cannot tell node B.
- Another client reads
xfrom node B.
Node B has only two choices:
- Answer with the old value. The system stays available but is no longer consistent.
- Refuse to answer (return an error or wait). The system stays consistent but is no longer available.
There is no third option. B cannot know about a write it never received.
"Pick two" is misleading
CAP is often drawn as a triangle with "pick any two". That suggests you could choose CA and skip partition tolerance. In a system spread across a network, you cannot: cables are cut, switches fail, and a long garbage collection pause is indistinguishable from a network failure. Partitions happen whether you tolerate them or not.
A single-node database is "CA" only in the trivial sense that it has no network between replicas, and when that one node is down, it is simply unavailable.
A better statement: when a partition occurs, choose between consistency and availability.
CP: consistency over availability
A CP system prefers refusing requests to giving wrong answers. Typically, only the side of the partition that holds a majority of nodes continues to accept writes. The minority side rejects requests until the network heals.
- Examples: etcd, ZooKeeper, Consul, HBase, and databases configured for strongly consistent replication. These rely on consensus protocols; see how Raft works.
- Use when being wrong is worse than being down: account balances, inventory that must not be oversold, locks and leader election, configuration that every node must agree on.
AP: availability over consistency
An AP system answers every request, even if nodes cannot talk to each other. Both sides of the partition accept writes. When the network heals, the system must reconcile the differences.
- Examples: Cassandra and Riak in their default configurations, DNS, and shopping carts built on Amazon's Dynamo design.
- Use when being down is worse than being briefly wrong: shopping carts, social feeds, likes and view counters, product catalogues, caches.
The reconciliation is where the complexity goes. See eventual consistency explained.
Two everyday examples
An ATM network. Suppose an ATM loses contact with the bank.
- A CP design refuses withdrawals: "service unavailable".
- An AP design allows withdrawals up to a small limit and reconciles later, accepting the risk of an overdraft.
Real banks often choose the second, because a limited overdraft costs less than angry customers. The business decides how much inconsistency is tolerable.
DNS. When a record changes, resolvers around the world keep serving the cached old value until it expires. DNS is always available and only eventually consistent. See how DNS works.
PACELC: what about when nothing is broken?
CAP only describes behaviour during a partition, which is rare. Most of the time the network is fine, and a different trade-off applies. The PACELC formulation adds it:
If there is a Partition, choose Availability or Consistency; Else, choose Latency or Consistency.
To guarantee consistency, a write must be confirmed by several replicas before it is acknowledged, and reads may need to check with a majority. That takes time. Skipping the coordination is faster but allows stale reads. This latency cost is paid on every request, so in practice it shapes system design more than partition behaviour does.
| System | During a partition | Otherwise |
|---|---|---|
| Cassandra, DynamoDB (default settings) | Availability | Low latency |
| Spanner, etcd, CockroachDB | Consistency | Consistency (at a latency cost) |
| MongoDB (default settings) | Consistency | Tunable |
The limits of CAP
Martin Kleppmann's article Please stop calling databases CP or AP explains why the labels are too crude:
- Consistency is a spectrum. Between linearizability and eventual consistency lie many useful models, such as snapshot isolation and causal consistency. The Jepsen consistency map lays them out.
- Availability is a spectrum. No real system is available 100% of the time.
- Many systems are tunable per request. In Cassandra you choose, for each query, how many replicas must respond.
- Some systems are neither. A leader-follower database with asynchronous replication can lose availability and return stale data during a partition.
- It says nothing about performance, durability or cost.
Treat CAP as a way of starting the conversation, not as a classification scheme.
How to use it in practice
- Ask, for each piece of data: what is the cost of serving a stale or conflicting value? What is the cost of refusing the request?
- Different data in the same product can make different choices. The cart can be AP while payment is CP.
- Check what your database actually does under failure, with your configuration. Defaults vary, and so do the results of independent testing.
- Remember the latency trade-off, which affects every request.
For how data is copied between nodes in the first place, see how database replication works.
Frequently asked questions
What does CAP stand for?
Consistency, Availability and Partition tolerance.
Can a system be CA?
Not a distributed one. Network partitions cannot be ruled out, so you must decide how to behave when one happens.
Is MongoDB CP or AP?
With default settings it behaves as CP: during a partition, only the side with a majority elects a primary and accepts writes. Its read and write settings can relax this.
Does CAP mean NoSQL databases are inconsistent?
No. Many offer strong consistency as an option, and several distributed SQL databases provide it by default. CAP describes a trade-off, not a property of any category.
Conclusion
The CAP theorem captures one unavoidable fact: when nodes cannot communicate, a system must either stop answering or risk answering wrongly. Which is worse depends on the data. Use CAP to frame that decision, PACELC to remember the everyday latency cost, and your database's real documented behaviour to make the final call.
Related articles
- What Is Eventual Consistency and When Is It Acceptable?
- How Database Replication Works
- How Consensus Algorithms Like Raft Keep Servers in Agreement
- Why Distributed Systems Are So Hard
