Command Palette

Search for a command to run...

Hectal
PHASE 10Advanced ~8 min· topic 2 of 5

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.

0/5 · 0%

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

  1. 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.

  2. 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).

  3. 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.

  4. 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.

  5. 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

quorum.txttext
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 writes
cap.mermaiddiagram
Rendering diagram…

Interview 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

01

Why is "CA" not a real option for distributed databases?

Practice

01

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.