CAP Theorem
Consistency, Availability, Partition Tolerance
- In Turkish
- CAP Teoremi
- Pronunciation
- KAP THEER-um
In short
The CAP theorem says that if a network failure splits a distributed database, the system must choose between consistency and availability; it can't have both.
What is the CAP theorem?
The CAP theorem describes a fundamental trade-off in distributed data systems, meaning databases that store data on several networked machines. It names three properties: consistency (every read returns the most recent write or an error), availability (every request to a working node gets a non-error response), and partition tolerance (the system keeps working even when network failures cut some nodes off from others). It was proposed by Eric Brewer in 2000 and proven by Seth Gilbert and Nancy Lynch in 2002.
It is often summarized as 'pick two of three', but that is misleading. Network partitions can't be prevented in a real distributed system, so partition tolerance is not optional, and the actual choice is what to do while a partition lasts. A CP system refuses or delays some requests to avoid returning stale data, while an AP system keeps answering but may return out-of-date data until the network heals.
Imagine two bank branches that lose their phone line to each other. Either they stop allowing withdrawals until the line is back, choosing consistency, or they keep serving customers and reconcile balances later, choosing availability and risking an overdraft. Systems that handle payments or inventory often lean CP, while shopping carts, social feeds, and DNS usually lean AP and rely on eventual consistency, where all copies converge once communication is restored.
Two confusions are worth knowing. The C in CAP means linearizability, a strict guarantee that every node returns up-to-date data, which is different from the C in ACID, where consistency means the data follows the database's rules and constraints. And CAP only describes behavior during partitions; the PACELC theorem extends it by noting that even when the network is healthy, systems trade latency against consistency.
At a glance
Key takeaways
- CAP stands for consistency, availability, and partition tolerance.
- During a network partition, a distributed system must choose consistency or availability.
- Partition tolerance is mandatory in practice, so the real choice is CP or AP.
- The C in CAP is not the same as the C in ACID.
- PACELC extends CAP to the latency-versus-consistency trade-off in normal operation.
Example
-- In cqlsh, the Cassandra shell, the trade-off can be tuned per request
-- Favors availability: any single replica may answer, possibly with stale data
CONSISTENCY ONE;
SELECT balance FROM accounts WHERE id = 42;
-- Favors consistency: a majority of replicas must respond, or the read fails
CONSISTENCY QUORUM;
SELECT balance FROM accounts WHERE id = 42;Readers ask
Which is more important, consistency or availability?
It depends on the data. Bank balances and inventory counts usually need consistency, while social feeds, view counts, and shopping carts can tolerate brief staleness in exchange for staying available.
Does the CAP theorem apply to a single-server database?
Not really. CAP applies to data replicated across multiple networked nodes; a single server has no partition between copies to worry about, although it can still simply go down.
What is eventual consistency?
Eventual consistency means that if no new updates are made, all copies of the data will eventually become identical. It is the typical guarantee of AP systems, which stay available during partitions and reconcile differences afterward.
See also
- Database ReplicationDatabases, p. 10Database replication is the continuous copying of data from one database server to others, so several servers hold the same data for reliability and scale.
- ShardingDatabases, p. 39Sharding is a way of scaling a database by splitting its data across several servers, called shards, so each one stores and handles only part of the total.
- NoSQLDatabases, p. 29NoSQL is a family of databases that store data in models other than relational tables, such as documents, key-value pairs, wide columns, or graphs.
- TransactionDatabases, p. 47A transaction is a group of database operations that succeed or fail as a single unit, so the data is never left in a half-finished, inconsistent state.
- DatabaseDatabases, p. 6A database is an organized collection of data stored on a computer, managed by software that lets applications save, search, and update it efficiently.
- MicroservicesSoftware Architecture, p. 27Microservices are an architectural style where an application is split into small, independently deployable services that communicate over a network.
- Distributed SystemSoftware Architecture, p. 14A distributed system is a set of computers that work together over a network and appear to their users as a single system.
- Consensus AlgorithmSoftware Architecture, p. 8A consensus algorithm lets a group of machines agree on a single value or an ordered log of decisions, even when some of them crash or messages are lost.
Sources
Spotted a mistake or something missing on this page?Suggest an edit