System Design Guide

Fundamentals · 13

CAP theorem, PACELC and consistency models

What CAP really says, why it's about network partitions, how PACELC adds the everyday latency trade-off, and the consistency models between "strong" and "eventual".

3 min read · 6 flashcards

CAP, stated precisely

In a distributed data store you’d like all three of:

  • Consistency ©: every read sees the most recent write. Formally, linearizability: the system behaves like one up-to-date copy.
  • Availability (A): every request to a working node gets a (non-error) response.
  • Partition tolerance (P): the system keeps working when the network drops messages between nodes.

The CAP theorem says that during a network partition you can’t have both C and A. Partitions happen in any real network, so P is not optional, and the meaningful question is: when a partition happens, do you refuse some requests (CP) or serve possibly stale data (AP)?

Client AClient BNode 1x = 2 (new)Node 2x = 1 (old)write x = 2read x ?network partition
A partition: the two sides can't talk. Answer anyway (AP) or refuse (CP)?

Client B asks node 2 for x while node 2 can’t hear from node 1:

  • A CP system makes node 2 return an error or wait, staying correct but unavailable.
  • An AP system makes node 2 return x = 1, staying available but stale. The nodes reconcile when the partition heals.

CP vs AP in practice

CP: choose consistency

  • Minority side of a partition stops accepting writes
  • Never returns stale or conflicting data
  • Use for: payments, inventory, bookings, locks, leader election
  • Examples: ZooKeeper, etcd, Spanner, HBase, a single-leader SQL setup

AP: choose availability

  • Every side keeps serving; replicas converge later
  • May return stale data or create conflicts to resolve
  • Use for: social feeds, likes, shopping carts, DNS, product catalogue
  • Examples: Cassandra, DynamoDB (default), CouchDB, Riak

Many systems are tunable: Cassandra and DynamoDB let you choose per query (quorum reads for consistency, single-replica reads for speed).

PACELC: the everyday trade-off

Partitions are rare, but replication is constant. PACELC extends CAP:

If there’s a Partition, choose A or C; Else, choose Latency or Consistency.

Even on a healthy network, waiting for replicas to confirm (consistency) costs latency, and not waiting (low latency) risks stale reads. Examples: DynamoDB and Cassandra are PA/EL (available, low latency); Spanner and most SQL setups are PC/EC (consistent, paying latency).

Consistency models, strongest to weakest

Model Guarantee Typical use
Strong / linearizable Every read sees the latest write; behaves like one copy Balances, locks, unique usernames
Sequential Everyone sees operations in the same order (not necessarily real-time) Replicated logs
Causal Cause comes before effect (a reply never appears before its question) Comments, chat
Read-your-writes You always see your own updates Profile edits, posts
Monotonic reads You never see data go “back in time” Any feed
Eventual Replicas converge if writes stop Likes, view counts, DNS

Test yourself

Answer in your head, then click a card to check. All cards are in the Anki deck.