Topic 10.2
CAP, PACELC and Quorums
In one line
CAP says that during a network partition a system must choose between consistency (refuse some requests) and availability (answer with possibly stale data). PACELC adds that even without partitions there's a trade between latency and consistency. Quorum systems tune this with N replicas, W write acks and R read replicas: R + W > N gives overlapping reads and writes.
Think of it like this
Two bank branches whose phone line is cut. Either they stop giving out cash until the line returns (consistent, not available) or they keep paying out and risk the account going overdrawn across both branches (available, not consistent).
Key ideas
- 01
Partitions aren't optional in distributed systems, so the real choice is CP or AP during a partition. CP: etcd, ZooKeeper, a PostgreSQL primary with synchronous replication, Spanner. AP: Cassandra and DynamoDB with eventual reads, DNS.
- 02
CAP's C is linearizability and its A is "every non-failed node answers". Most real systems are neither fully; they're configurable per operation (Cassandra consistency levels, DynamoDB strongly consistent reads, MongoDB read/write concerns).
- 03
PACELC: if Partition, choose A or C; Else, choose Latency or Consistency. Spanner is PC/EC (consistent, pays latency); Cassandra at ONE is PA/EL; DynamoDB default is PA/EL with optional strong reads.
- 04
Quorums: with N = 3, W = 2, R = 2, every read overlaps at least one replica with the latest write. W = 1, R = 1 is fastest and weakest. Quorums alone aren't linearizable without read repair and careful handling of concurrent and failed writes.
- 05
Leader-based vs leaderless: leader-based (PostgreSQL, MySQL, MongoDB, Kafka) orders writes through one node per shard; leaderless (Cassandra, Dynamo-style) lets any replica accept writes and reconciles via quorums, hinted handoff, read repair and anti-entropy.
Code & diagrams
N = replicas, W = acks needed for a write, R = replicas consulted on a read
N=3 W=2 R=2 -> R+W=4 > 3 : reads overlap the latest write, tolerates 1 node down
N=3 W=3 R=1 -> fast reads, writes fail if any replica is down
N=3 W=1 R=1 -> lowest latency, stale reads possible
Cassandra: CONSISTENCY LOCAL_QUORUM in each DC; EACH_QUORUM for cross-DC writesInterview problem
The problem
Classify and configure a shopping cart and a payment ledger
Using PACELC, decide the behaviour of (a) a shopping cart replicated across 3 regions and (b) a payment ledger, during partitions and in normal operation, and choose configurations.
When it breaks
Assuming QUORUM gives strong consistency
What you see
A failed write that reached only 1 of 3 replicas is later read by a quorum including that replica; read repair spreads a value the client was told had failed.
Fix & prevent
Use idempotent, versioned writes; understand that quorum systems aren't linearizable by default; use lightweight transactions (Paxos) for compare-and-set.
Explain it without notes
Why is "CA" not a real option for distributed databases?
Practice
For N = 5, list the (R, W) pairs that guarantee overlap and tolerate 2 failed nodes for reads.
Trade-offs
- ↔
CP protects invariants at the cost of availability; AP keeps serving at the cost of reconciliation; PACELC reminds you latency is the everyday price of consistency.
Done when you can
I can explain CAP and PACELC precisely and configure quorums for a requirement.