Lesson 04

CAP Theorem

Consistency vs availability vs partition tolerance — when to choose what.

Plays in the sticky player at the bottom of the page

Transcript

Lesson 4: CAP Theorem

A reference-style deep dive into consistency vs availability vs partition tolerance — when to choose what. Read this as an article, not a transcript. The accompanying audio is the spoken companion; the article below is the canonical written reference.

Audience. Engineers designing, building, or operating distributed systems who want a clear mental model rather than a checklist of tools.

Prerequisites. Working knowledge of HTTP, basic SQL, and the idea of running more than one server behind a load balancer.

Table of contents

  1. Why cap theorem matters
  2. The core mental model and where to start 3-12. In-depth sections below (see "Lesson body")

Lesson diagram

Lesson 4 diagram — CAP Theorem

Figure 1. The canonical cap theorem topology and control flow covered in this lesson.


1. Why CAP matters

CAP is the most-cited and most-misunderstood theorem in distributed systems. It says: during a network partition, you must choose between consistency and availability.

The catch is that partitions are inevitable. Network switches fail. Routers reboot. Fiber gets cut. The question is not "if" but "when" and "how often". CAP forces you to choose in advance: when a partition happens, do you serve (potentially stale) data, or do you refuse to serve?

Eric Brewer's original CAP formulation (2000) was informal. Seth Gilbert and Nancy Lynch's proof (2002) formalized it. The proof has been refined over the years; the core insight remains.

The common misconception is that you can pick CA (no partition tolerance). You cannot, in any real distributed system. Networks fail. So the real choice is CP or AP.

2. The formal statement

In any distributed data store, you can have at most two of:

  • Consistency — every read receives the most recent write or an error.
  • Availability — every request receives a (non-error) response, without the guarantee that it contains the most recent write.
  • Partition tolerance — the system continues to operate despite network partitions between nodes.

A "partition" is any message loss between nodes. Since message loss is inevitable in real networks, partition tolerance is required. So the real choice is CP or AP.

# CP system
def read(key):
    if not can_reach_majority():
        raise UnavailableError()    # refuse rather than return stale
    return quorum_read(key)

# AP system
def read(key):
    return any_replica_read(key)    # may return stale data

3. Consistency levels — the spectrum

CAP is binary, but consistency is a spectrum. The real-world options:

Linearizability (strongest)

Every operation appears to take effect atomically at some point between invocation and response. The system behaves as if there is a single copy of the data.

# Linearizable counter
class LinearizableCounter:
    def __init__(self):
        self.value = 0
    def increment(self):
        # Acquire lock; modify; release
        with self.lock:
            self.value += 1
            return self.value

Linearizability is expensive: requires consensus or coordination on every operation. Most production systems cannot afford it everywhere.

Sequential consistency

Operations appear in some global order consistent with each client's program order. Weaker than linearizability: two concurrent operations can be observed in either order.

Causal consistency

Operations that are causally related are seen in causal order; independent operations can be reordered.

# Causal: read your own writes
session.write(key, value)
session.read(key)   # MUST see the value just written

Per-session causal consistency is achievable cheaply (no global coordination).

Read-your-writes consistency

A client always sees its own writes. Achievable with sticky sessions to the master.

Monotonic reads consistency

A client never sees older data after seeing newer data.

Eventual consistency (weakest)

If no new writes happen, all replicas will eventually converge. The default for most distributed systems.

4. PACELC — extending CAP

CAP only describes partition behavior. PACELC extends it: even when there is no partition, you still have a trade-off between latency and consistency.

P → during a partition, choose A or C
E → else (no partition), choose L (latency) or C (consistency)

So a system is described by two letters: PC/EC, PA/EL, PC/EL, PA/EC.

SystemPACELCWhy
DynamoDBPA/ELTunable: per-request consistency
CassandraPA/ELTunable
HBasePC/ECStrong consistency via ZooKeeper
MongoDBPC/ECStrong consistency by default
CouchDBPA/ELEventual consistency

The PACELC framing is more useful than CAP because it captures the everyday trade-off (every request is a latency-vs-consistency choice) rather than the rare-event trade-off (partitions).

5. Real systems and their CAP choices

CP systems

  • HBase — relies on ZooKeeper; refuses writes when minority.
  • MongoDB — strong consistency by default; some configurations allow AP.
  • etcd / Consul — used for coordination; strong consistency required.
  • Redis Cluster — configurable; default is AP, but transactions make it CP for those operations.

AP systems

  • DynamoDB — always serves; tunable per-request consistency.
  • Cassandra — always serves; tunable per-query.
  • CouchDB — eventual consistency by design.
  • Riak — eventual consistency; tunable.

The choice in practice

The right choice depends on the workload:

  • Money — CP. Refuse to serve rather than return the wrong balance.
  • Social feed — AP. A stale feed for 5 seconds is fine.
  • Multiplayer game — AP with bounded staleness. The game must stay interactive.
  • DNS — AP. Stale DNS is fine; the cost of unavailability is higher.

6. PACELC in the real world

Every system makes a default PACELC choice; many let you override per request.

# Cassandra: per-query consistency
session.execute(query, consistency_level=ConsistencyLevel.QUORUM)
session.execute(query, consistency_level=ConsistencyLevel.ONE)   # faster, less consistent

# DynamoDB: per-request consistency
table.get_item(Key=..., ConsistentRead=True)   # strong
table.get_item(Key=..., ConsistentRead=False)  # eventual, faster

The right knob for a given request depends on its business impact. A friend list read can be eventual; an account balance read must be strong.

7. The "consistency window" of async replication

When you choose AP, you accept a consistency window: the worst-case staleness a read may see.

# For an AP system with async replication:
#   - Write returns to client at t=0
#   - Replication lag is p50=10ms, p99=500ms
#   - Worst case: a read at t=10s after the write sees the new value
#     → consistency window = replication lag

Monitoring the consistency window is essential:

# Prometheus: replication lag in seconds
max(replication_lag_seconds) by (cluster)

The alert should fire when the consistency window exceeds your business's tolerance.

8. CAP-violating claims

A claim that violates CAP usually means one of two things:

  1. The system is not actually distributed. A single-node database trivially provides both C and A because there is no partition to tolerate.
  2. The system makes the trade-off in a non-obvious way. It might serve during a partition but in a degraded mode (e.g. read-only). Or it might use a different consistency definition.

When someone claims a system violates CAP, ask: "during a partition, what happens to a write that cannot reach a quorum?" The answer reveals the actual CAP choice.

9. The CAP trade-off is not the only one

CAP describes the partition trade-off. PACELC adds the no-partition trade-off. But there are others:

  • Latency vs throughput — more consistency usually means more coordination, which means more latency per request, which means fewer requests per second.
  • Cost vs reliability — synchronous replication costs more than async (more machines, more network).
  • Operational complexity — running a CP system is harder than running an AP system; you must handle the unavailability cases.

The CAP trade-off is the most-cited but rarely the most-important one in practice.

10. Decision framework

Use this framework to choose C or A:

What is the cost of a stale read?
├── High (money, safety-critical)
│   └── Use CP. Refuse to serve rather than be wrong.
└── Low (social, content, cache)
    └── Use AP. Stay available; accept staleness.

What is the cost of unavailability?
├── High (DNS, login, checkout)
│   └── Bias toward AP; degrade gracefully rather than refuse.
└── Low (analytics, batch)
    └── Bias toward CP; correctness matters more than speed.

When the cost of stale read is high AND the cost of unavailability is high, you have a hard problem. Solutions include: external consistency (e.g. version vectors), or business-level reconciliation.

11. The future: beyond CAP

CAP is a 2000-vintage theorem. Modern systems are exploring:

  • Conflict-free replicated data types (CRDTs) — merge without coordination.
  • Operational transformation — like Google Docs; merge concurrent edits logically.
  • Calvin — deterministic transaction ordering; no locks needed.
  • FoundationDB — combines a consensus-ordered log with optimistic concurrency.

These systems try to give you the best of both: availability and consistency, by being smart about which conflicts can be resolved automatically and which require coordination.

12. Key takeaways

  • CAP says: during a partition, choose C or A. P is not optional in real networks.
  • PACELC extends CAP: even without partitions, you trade latency for consistency.
  • CP systems refuse writes during a partition; AP systems serve stale data.
  • The right choice depends on the workload's cost of staleness and unavailability.
  • Tunable consistency (Dynamo, Cassandra) lets you pick per request.
  • Beyond CAP: CRDTs, OT, and Calvin explore new trade-offs.

Appendix: terms

  • Partition — any loss of communication between nodes.
  • Linearizability — strongest single-object consistency: every op appears atomic.
  • Sequential consistency — operations in some global order consistent with each client's program order.
  • Causal consistency — causally-related ops in causal order; independent ops can be reordered.
  • Eventual consistency — replicas converge given enough time without new writes.
  • PACELC — extension of CAP that adds the no-partition trade-off (latency vs consistency).

Appendix: source dialogue excerpt

The audio for this lesson was synthesized from the following Cantonese dialogue (verbatim, not translated):

  • M: 各位同學早晨, 我係子謙。歡迎收聽系統架構課程第四課。今日嘅主題係 CAP theorem。…
  • F: 大家好, 我係曉晴。CAP 係 distributed system 嘅 fundamental theory, 由 Eric Brewer 喺 2000 年提出, 2002 年由 Gilbert 同 Lynch 證明。今日我哋會拆解 CAP 嘅三個 property, 同埋佢哋之間嘅 tension, 同實際 dis…
  • M: 首先講解基本概念。CAP 嘅 C 係 Consistency, 即係所有 node 喺同一時間睇到同一份 data。A 係 Availability, 即係每個 request 都會收到 response, 唔係 error。P 係 Partition tolerance, 即係 network 兩個 node 之間嘅…
  • F: CAP 嘅 theorem 講, 喺一個 asynchronous network 入面, 你唔可以同時 satisfy 三個 property, 最多只能 satisfy 兩個。即係當 network partition 發生嘅時候, 你必須喺 consistency 同 availability 之間揀一個。Par…
  • M: 好, 第一個 important nuance 係 CAP 講嘅係 network partition 期間嘅 behavior, 唔係 steady state。正常 operation 嘅時候, 即係冇 partition, 系統可以同時 consistency 同 availability。CAP 嘅真正 imp…
  • F: CAP 嘅 CP model 即係 consistency 加 partition tolerance。例如 HBase, ZooKeeper, 嗰啲系統。當 partition 發生嘅時候, 系統會 reject 一部分 request, 確保 consistency。例如 ZooKeeper 喺 partition…

Full dialogue contains 27 segments; see script_raw.json in the source folder.

Lesson quiz · 30 questions

Question 1 of 30Answered 0 / 30
Question 1 of 30

CAP theorem says during a network partition you must choose between...

Pick an answer to lock it in. We'll tell you immediately whether you got it right and show an explanation. Then press Enter or click Next to continue.

Shortcuts:ABCDpick answer on current questionEntergo to next unanswered
30 unanswered