The CAP Theorem
Analyze Eric Brewer's CAP theorem: Consistency vs Availability during a Network Partition. Understand why "CA" systems do not exist in reality.
CAP Theorem Triangle & The Partition Reality ⚖️
Network partitions are physical inevitabilities; systems choose CP (Consistency) or AP (Availability).
01.1. Deconstructing Brewer's CAP Guarantees
Formulated by Eric Brewer in 2000 and mathematically proven by Seth Gilbert and Nancy Lynch in 2002, the CAP Theorem states that a distributed data store can simultaneously provide at most two of the following three guarantees:
1. Consistency (Linearizability / Single-Copy Serializability)
Every read request receives the most recent write or an error. To a client, the cluster behaves as if there is only a single atomic, instantaneous copy of the data item in existence, regardless of how many replica nodes exist.
2. Availability (Liveness)
Every non-failing node must return a non-error response for every received read or write request (without guarantee that the response contains the absolute newest write). Note: In CAP proof terms, returning an HTTP 500 error or timing out means the system is unavailable.
3. Partition Tolerance (Fault Tolerance)
The system continues to operate despite an arbitrary number of messages being dropped, delayed, or partitioned by the underlying physical network connecting the nodes.
02.2. The Partition Reality: Why "CA" Does Not Exist in Distributed Systems
Many introductory texts present CAP as a 3-way Venn diagram where you can pick any 2 of 3 (CA, CP, or AP). In real-world distributed systems, this is a dangerous misconception.
Network partitions are not an optional design feature you can disable; they are a physical certainty of networking infrastructure (switches fail, optical transceivers degrade, DNS records desynchronize).
Therefore, Partition Tolerance (P) is mandatory. The actual CAP theorem formulation is:
In the presence of a network partition, you must choose between:
- CP (Consistency / Reject Availability):
- When a network split occurs between Datacenter A (US) and Datacenter B (EU), nodes in Datacenter B cannot reach the quorum leader in Datacenter A.
- To prevent serving stale data or recording conflicting writes, Datacenter B rejects incoming writes and returns errors.
- Result: 100% data correctness is preserved, but service uptime is sacrificed.
- AP (Availability / Sacrifice Consistency):
- Nodes in Datacenter B continue accepting client writes and serving local reads using whatever data they currently hold.
- When the network partition eventually heals, background synchronization or Conflict-Free Replicated Data Types (CRDTs) reconcile conflicting updates.
- Result: 100% uptime SLA is maintained, but clients experience stale reads and temporary inconsistencies.
03.3. Real-World Database Taxonomy under CAP
| Database Category | CAP Classification | Primary Use Cases | Behavior During Partition |
|---|---|---|---|
| Google Cloud Spanner / CockroachDB | CP | Financial ledgers, ACID checkouts, banking. | Rejects transactions on partitioned minority nodes; guarantees zero stale reads. |
| Apache ZooKeeper / etcd | CP | Distributed lock managers, leader election, metadata. | Halts operations if a leader cannot maintain a quorum majority (> N/2). |
| Apache Cassandra / ScyllaDB | AP (Tunable) | Telemetry, IoT sensor data, user activity feeds. | Accepts writes on any available node; synchronizes via hinted handoff and read repair. |
| Amazon DynamoDB | AP (Default) | E-commerce shopping carts, session stores. | Prioritizes sub-10ms response availability over immediate linearizability. |
⚖️Architectural Trade-offs & Production Realities
Architectural Advantages
- Provides a rigorous, universally understood theoretical model for classifying distributed data stores
- Forces system architects to explicitly make trade-offs between consistency correctness and availability SLAs
Trade-offs & Constraints
- Oversimplifies real-world engineering by ignoring latency during normal (non-partitioned) operations
- The binary definition of "Availability" (100% non-error response) does not reflect modern partial degradation techniques
Amazon famously chose AP for the Dynamo shopping cart: A customer must always be able to add an item to their cart even during cross-datacenter fiber partitions, resolving split-brain cart versions later during checkout.
🎯 Staff+ Engineering Takeaways
- Network partitions are physical facts; trade-offs only occur when partitions strike.
- You cannot "choose CA" in a distributed network—P is mandatory.
- CP prioritizes correctness over uptime (rejects writes if quorum is lost).
- AP prioritizes uptime over instant consistency (accepts writes and resolves conflicts later).
Topic Knowledge Assessment 🧠
Step through 1 scenario question to test your staff-level grasp.
Why can a distributed database not choose "CA" (Consistency + Availability without Partition Tolerance)?
How clear and staff-actionable was this system breakdown?