The CAP Theorem

Suppose the merchant changed the price of the blue mug just before the connection between the primary and the replica broke. The primary recorded the new price, but the change had not reached the replica. Now a shopper browses the mug shop’s catalog, and the read goes to the replica. The replica can answer with the old price it has. Or it can refuse to answer until it can reach the primary again.

Every replicated database faces this choice during a partition. The CAP theorem describes it.

Consistency, availability, and partition tolerance

The CAP theorem describes a limit on the guarantees a distributed data system can provide. Such a system stores its data on several servers that communicate over a network. Shopend’s replicated database is one. The theorem is named after three properties:

  • Consistency ©: reads and writes behave as though the data were stored on one up-to-date server. A read returns the result of the latest write that completed before the read began. For Shopend, once the new price is recorded, every later read of the product returns the new price. It does not matter which replica answers the read.
  • Availability (A): every request that a running database server receives gets a response. By response, I mean the request completes: a read returns the data, and a write is recorded. An error does not count. Neither does holding the request without answering until the server can reach the other servers again. The data in the response may not be up to date (that is what consistency is about). For Shopend, the replica must answer the catalog read, and the primary must complete the writes it receives.
  • Partition tolerance (P): the system keeps operating when communication between groups of servers is lost. The network partition described earlier is an example of that kind of loss.

Availability as uptime

We have used the word availability before, in a different sense. Shopend keeps the shop API’s availability requirement: during each calendar month, clients must be able to use the API at least 99.9 percent of the time. That is availability as uptime.

Availability in CAP is narrower. It is a requirement on each request rather than on a share of time: during a partition, every request that a running server receives must get a response. In our example, the majority rule makes the primary refuse writes while the partition lasts. If the partition lasts a few minutes, Shopend may still meet its 99.9 percent availability requirement for the month. But in CAP’s sense, the database was not available for writes during those minutes.

Two meanings of consistency

We used the word consistency earlier, for ACID transactions. In Dining Dollars, a purchase cannot make a student’s balance negative, and the database checks that rule on every write. In ACID, consistency means a transaction takes the database from a state where every rule on the data holds to another state where every rule holds.

In CAP, consistency is about what clients observe when the data is replicated.

The two are independent. Suppose the last blue mug sold just before the connection broke. The primary has a quantity of 0, but the replica may still have a quantity of 1. A read from the replica returns an old value, so the system is not consistent in CAP’s sense. But the replica’s data still follows the rules on the data. No stock quantity is below zero. In ACID terms, the replica’s data is consistent.

Consistency or availability during a partition

Go back to the replica during the partition. It receives the catalog read, and the new price exists only on the primary.

If it answers with the product it has, the response shows the old price. The system stays available and gives up consistency.

If it refuses, or waits until it can reach the primary, it never returns the old price. The system stays consistent and gives up availability for that read.

This is the CAP theorem: a distributed data system cannot guarantee both consistency and availability during a network partition. To preserve consistency, it may have to stop completing requests that need communication across the partition. To preserve availability, it must let each server complete requests with the data it has. The responses may then disagree or be out of date.

Eric Brewer proposed this limit as a conjecture in a 2000 conference talk. Seth Gilbert and Nancy Lynch proved it in 2002. Since then it has been called the CAP theorem.