# The CAP Theorem

> In a distributed system, a network partition forces a choice between consistency and availability - what CAP actually claims, why the choice is unavoidable, and what it looks like in real databases.


---

# The CAP Theorem

Every distributed database makes a promise, and eventually breaks part of it - on purpose, because the network forced it to. The CAP theorem is the formal statement of which part breaks, and when. It sounds abstract until the actual moment it applies: a cable gets cut, two data centers can't talk, and every system behind that cable must decide, right now, whether to keep answering or keep agreeing - it can't do both.

## How to read this

Read it in order. Phase 1 defines the three letters precisely - most confusion about CAP comes from vague definitions of "consistency" and "availability." Phase 2 walks through the actual argument for why you can't have all three, using a concrete network partition, and Phase 3 grounds it in real databases you've probably used and clears up the most common misreading of the theorem.

## The phases

1. [The three letters](01-the-three-letters.md) - consistency, availability, and partition tolerance, defined precisely.
2. [Why you can't have all three](02-why-you-cant-have-all-three.md) - walking through a real partition and the forced choice it creates.
3. [What this looks like in real databases](03-what-this-looks-like-in-real-databases.md) - CP systems, AP systems, and the most common misreading of CAP.


---

# The three letters

CAP stands for Consistency, Availability, and Partition tolerance. Each word already has an everyday meaning that will mislead you slightly, so before anything else, here's what each one means specifically in this context - a distributed system, meaning data that's copied across more than one machine.

## C - Consistency

**In this context, consistency means every node in the system agrees on the current value.** If you write a new value to the system and then immediately read it back - from any node - you get that new value, not an old one. There is one true answer, and every copy reflects it.

```text
Node A: balance = $500
Node B: balance = $500
Node C: balance = $500

-> a client writes balance = $400 to Node A
-> consistency means A, B, and C ALL now say $400, before any client
   is allowed to read $500 again from any of them
```

*What just happened:* consistency here is entirely about agreement across copies of data, not about a single database enforcing rules like "balances can't go negative." That second meaning is a real thing too - it's the "C" in ACID - but it's a different concept that happens to share a letter. In CAP, consistency is purely: does everyone see the same, latest value, right now.

## A - Availability

**Availability means every request that reaches a working node gets a response - not an error, not a timeout, an actual answer.** It says nothing about whether that answer is the newest value. It only promises that the system answers.

```text
Client sends: GET balance

Available system's promise:
  -> you WILL get a response (even if it might be slightly stale)

NOT the promise:
  -> the response is guaranteed to be the very latest value
```

*What just happened:* availability is a promise about responsiveness, full stop. A system can be perfectly available and still hand you data that's a few seconds out of date. Those are two separate axes, and CAP is specifically about the tension between them - which is exactly why conflating "available" with "correct" is the fastest way to misread this theorem.

## P - Partition tolerance

**A partition is when the network between nodes breaks** - a cable gets cut, a data center loses connectivity, a router drops packets - and some nodes can no longer talk to others, even though every individual node is still running fine. **Partition tolerance means the system keeps operating in some form despite that split**, rather than halting entirely.

```text
Before the partition:
  [Node A] <---> [Node B] <---> [Node C]     (all connected, one system)

During a partition:
  [Node A] <---> [Node B]      X      [Node C]
                          (network split - C is cut off from A and B)
```

*What just happened:* Node C hasn't crashed - it's still running, still has data, and can still answer requests from clients that reach it. What it can't do is talk to A and B to confirm it has the latest value or hear about new writes. Partition tolerance is about what happens next: does the system keep working in this split state, or stop entirely until the network heals?

> Consistency is about agreement. Availability is about responsiveness. Partition tolerance is about surviving a broken network instead of giving up.

## Why partition tolerance isn't really optional

It's tempting to read CAP as "pick any two of three" - as if you could skip partition tolerance and keep both consistency and availability. That's not a real option: every distributed system's nodes talk over a network, and networks fail - cables get cut, switches misbehave, cloud regions have outages. A system that isn't partition-tolerant hasn't "chosen CA"; it just stops working correctly the moment a partition happens, which is a matter of when, not if.

```text
Not partition-tolerant  ->  breaks (in some way) the moment a real partition occurs
Partition-tolerant       ->  keeps running during a partition, but must choose:
                              stay consistent, or stay available. Not both.
```

*What just happened:* this reframes the whole theorem - it isn't really "pick two of three" as a free menu. Partition tolerance is a fact of life for any real distributed system; the actual decision you get to make is what to do *during* a partition: favor consistency, or favor availability. Phase 2 walks through exactly why you can't have both once the network is actually split.


---

# Why you can't have all three

Phase 1 defined the three letters. This phase is the actual theorem: walk through one concrete partition, step by step, and watch consistency and availability become mutually exclusive in real time. This isn't an abstract impossibility proof - it's a forced decision that a real piece of software has to make, right now, with an actual request waiting on it.

## The setup

Two nodes, A and B, both holding a copy of the same value, replicated so either one can serve reads and writes.

```text
Node A: inventory_count = 10
Node B: inventory_count = 10

Client 1 is connected to Node A.
Client 2 is connected to Node B.
```

*What just happened:* nothing unusual yet - this is just ordinary replication, working exactly as intended. Both nodes agree, and a client can talk to either one and get the same, correct answer.

## The partition happens

The network link between A and B goes down. Both nodes are still running. Both can still serve the clients connected to them. They can't talk to each other anymore.

```text
Node A: inventory_count = 10        X        Node B: inventory_count = 10
              (Client 1)                            (Client 2)
                        (network link is down)
```

*What just happened:* this is the partition from Phase 1, made concrete. Neither node crashed. Neither client's connection dropped. The only thing broken is A and B's ability to synchronize with each other.

## Client 1 makes a write

While the partition is ongoing, Client 1 sells the last unit and writes `inventory_count = 9` to Node A.

```text
Node A: inventory_count = 9   <- just updated
Node B: inventory_count = 10  <- has no idea; can't hear from A
```

*What just happened:* Node A applied the write locally - it has no other choice, since it can't check with B first. Node B still thinks the count is 10, and has no way to find out otherwise, because the link that would carry that update is down.

## Client 2 makes a request - and now the system must choose

At this exact moment, Client 2 asks Node B for the current inventory count. Node B has to answer *something*, and every option it has falls into one of two camps.

```text
Option 1 - favor consistency:
  Node B says: "I can't confirm I have the latest value, since I can't reach A.
                I will not answer." -> returns an error / times out
  Result: the system just gave up AVAILABILITY to protect CONSISTENCY.

Option 2 - favor availability:
  Node B says: "Here's what I have: inventory_count = 10."
  Result: Client 2 just got a STALE, WRONG answer (it's actually 9).
          The system gave up CONSISTENCY to protect AVAILABILITY.
```

*What just happened:* there is no third option. Node B cannot both answer *and* guarantee correctness, since confirming correctness requires reaching A, and A is unreachable - every possible response falls cleanly into one of the two buckets above. This is the theorem itself, not a metaphor for it: an actual node, facing an actual request, during an actual partition, with only two shapes of response available.

> During a partition, "answer, but maybe wrong" and "refuse to answer" are the only two moves left on the board. CAP is the statement that there is no third move.

## Why this only bites during the partition

Before the partition and after it heals, A and B can talk, and the system can be both consistent and available at once - there's no tension because there's no obstacle to synchronizing. The forced choice is strictly a partition-time phenomenon.

```text
No partition       -> consistency AND availability, simultaneously, no conflict
During a partition -> pick one, for as long as the partition lasts
Partition heals     -> back to both, once the nodes resync
```

*What just happened:* this is worth sitting with, because it's the detail most summaries skip, and it sets up the correction in Phase 3 - CAP doesn't say a system is "a CP system" or "an AP system" as some permanent identity. It says: *when a partition happens, which way does this particular system lean.* The rest of the time, the theorem has nothing to say at all.

## What "choosing" looks like in practice

Real systems don't flip a switch mid-outage - the choice is a design decision, baked into the software in advance: what should a node do if it can't reach its peers? Some default to Option 1: refuse, or elect a leader and let only the leader answer. Others default to Option 2: always respond, and reconcile the disagreement once the partition heals - Phase 3 looks at real databases that made each choice, and clears up the most common misunderstanding about what CAP actually claims.

```quiz
[
  {
    "q": "During a network partition, Node B cannot reach Node A to confirm it has the latest value, but a client asks Node B for data anyway. What are Node B's only two options, according to CAP?",
    "choices": [
      "Answer with a value that might be stale, or refuse to answer at all",
      "Answer correctly no matter what, or shut down permanently",
      "Ask the client to wait until the network heals, with no other option",
      "Automatically fail over to a third node that always has the answer"
    ],
    "answer": 0,
    "explain": "Node B can't verify correctness without reaching A, so it can only answer (risking staleness, favoring availability) or refuse (favoring consistency). There's no third move once the partition is real."
  },
  {
    "q": "When does the tension between consistency and availability actually apply, according to this phase?",
    "choices": [
      "At all times, permanently, for any distributed system",
      "Only during an actual network partition - before and after, both can hold at once",
      "Only when the system has more than 10 nodes",
      "Only during scheduled maintenance windows"
    ],
    "answer": 1,
    "explain": "Outside a partition, nodes can synchronize freely and there's no obstacle to being both consistent and available. The forced choice is specifically a partition-time phenomenon."
  },
  {
    "q": "In the walkthrough, why can't Node B return the correct, up-to-date value while the partition is ongoing?",
    "choices": [
      "Node B's disk is full",
      "Node B doesn't support writes",
      "Confirming the latest value would require reaching Node A, which is unreachable during the partition",
      "The client's request timed out before Node B could respond"
    ],
    "answer": 2,
    "explain": "Node B's only copy of the truth is its own last-known value. Verifying it's still current requires the very connection to A that the partition has severed."
  }
]
```

Watch it animated: [the CAP theorem](/explainers/CAPTheorem.dc.html)


---

# What this looks like in real databases

Phase 2 showed the forced choice in the abstract: during a partition, a node either answers with a possibly-stale value or refuses to answer. Real databases make that choice deliberately, as a design decision baked in ahead of time. Systems are often labeled "CP" or "AP" as shorthand for which way they lean.

## CP systems: consistency over availability

A **CP system** is built to refuse rather than risk answering wrong. If it can't confirm it has the current, agreed-upon value, it will return an error or block until it can.

**A single-leader relational database during failover.** Most traditional RDBMS setups (a primary with one or more replicas) route all writes through one leader node. If that leader becomes unreachable - partitioned off from the rest of the cluster - the system doesn't let just any replica start accepting writes with unconfirmed data; it typically pauses writes (or serves only reads, clearly marked as such) until a new leader is safely elected or the original one rejoins. The gap in write availability is the price paid to avoid two nodes disagreeing about the current value.

**ZooKeeper and etcd.** These are coordination services, specifically built to be the single source of truth other systems rely on for things like leader election and distributed locks. If they answered with stale data during a partition, every system depending on them for correctness would inherit that corruption. So they're built explicitly CP: during a partition, a node that can't confirm it's part of the majority (the quorum) will refuse to serve requests rather than risk giving an answer that disagrees with the rest of the cluster.

```text
CP system during a partition:
  "I can't confirm this is current." -> refuses / errors / blocks
  Cost: some requests fail during the partition
  Benefit: no client ever sees two different "correct" answers
```

## AP systems: availability over consistency

An **AP system** is built to always answer, even if that answer might be a few moments out of date, and to reconcile any disagreements after the partition heals.

**DynamoDB and Cassandra-style eventual consistency.** These systems are designed so that any reachable node answers a request, partition or not. If Node A and Node B both accepted writes to the same record while cut off from each other, the system doesn't block - it lets both writes happen, and resolves the conflict later using a defined strategy (last-write-wins, version vectors, or an application-supplied merge rule). The read you get back is described as **eventually consistent**: it will converge to the same value everywhere, but not necessarily at the instant you ask.

```text
AP system during a partition:
  "Here's what I have." -> answers immediately, even if stale
  Cost: two nodes might briefly disagree about the current value
  Benefit: no client is ever left without an answer
```

*What just happened, comparing the two:* neither choice is a bug - both are correct engineering responses to the same physical fact (the network split), aimed at different priorities. A payment ledger and a coordination service both lean CP because a wrong answer is worse than no answer. A shopping cart or a social media feed usually leans AP because a few-seconds-stale "add to cart" is a far smaller problem than the entire cart page going down.

## The most common misreading of CAP

Here's the correction that matters most: **CAP describes what a system does *during* a partition - it is not a permanent label for what a system is at all times.** A database described as "AP" isn't somehow less consistent on an ordinary Tuesday with no network issues; outside of a partition, most systems are both consistent and available simultaneously, exactly as Phase 2 showed. The CP/AP label only tells you which guarantee gets dropped in the specific, temporary window when nodes can't talk to each other.

```text
Wrong reading:  "System X is AP, so it's generally less consistent."
Right reading:   "System X is AP, meaning: IF a partition happens,
                  it will keep answering requests rather than refuse them,
                  and consistency is the thing that temporarily gives."
```

*What just happened:* this distinction is exactly why CAP gets misquoted so often in casual conversation. People use "CP" and "AP" as if they were permanent personality traits of a database, when the theorem only ever makes a claim about partition-time behavior. Most of the time, for most systems, the theorem doesn't apply at all - there's no partition, so there's nothing to trade off.

## The deeper refinement: PACELC

CAP's silence about the *non-partition* case is exactly what a later refinement, called **PACELC**, was designed to address: **if there's a Partition, choose between Availability and Consistency (the CAP part) - Else (when the network is fine), choose between Latency and Consistency.** The insight is that even without a partition, there's often still a tradeoff between responding fast and confirming full consistency across replicas first. It's worth knowing the name exists - it's the deeper, more complete version of the same question CAP raises.

## Recap

1. **Consistency** means every node agrees on the latest value; **availability** means every request gets an answer; **partition tolerance** means the system keeps functioning despite a network split - and it's effectively mandatory for any real distributed system.
2. During an actual partition, a node facing a request can only answer (risking staleness) or refuse (risking availability) - there's no third option, which is the theorem itself, not a metaphor for it.
3. **CP systems** (a failing-over RDBMS, ZooKeeper, etcd) refuse rather than risk disagreement. **AP systems** (DynamoDB, Cassandra) always answer and reconcile later.
4. The most common misreading is treating CP/AP as a permanent label. CAP only describes behavior *during* a partition - outside of one, a system is typically both consistent and available.
5. PACELC is the deeper refinement: even without a partition, there's often still a latency-versus-consistency tradeoff to make.
