CAP Theorem, PACELC & Distributed Consensus Mechanics
Theoretical foundations of distributed systems: Brewer's CAP Theorem, Abadi's PACELC extension, Paxos and Raft consensus algorithms, and event-driven architectures.
01 — Notebook Information & Scope
“A distributed system is one in which the failure of a computer you didn’t even know existed can render your own computer unusable.” — Leslie Lamport. Network partitions are not an anomaly; they are an inevitable physical reality of distributed infrastructure.
- Domain: Technology & Systems Architecture
- Subject: Distributed Systems
- Core Reference Model: Eric Brewer’s CAP Theorem, Daniel Abadi’s PACELC, and Raft Consensus
02 — Brewer’s CAP Theorem & PACELC Extension
In any asynchronous networked system, three properties cannot be simultaneously guaranteed during a network partition:
- Consistency (C): Every read receives the most recent write or an error (linearizable consistency).
- Availability (A): Every non-failing node returns a non-error response for every received request (without guarantee of containing the latest write).
- Partition Tolerance (P): The system continues to operate despite an arbitrary number of messages being dropped or delayed by the network between nodes.
Because physical networks will inevitably partition ($P$), a distributed database must choose between Consistency (CP) (refusing stale reads or blocking writes) or Availability (AP) (returning potentially stale data).
Abadi’s PACELC Extension
Daniel Abadi recognized that CAP only describes system trade-offs during a partition:
- If Partition ($P$): Choose between Availability ($A$) or Consistency ($C$).
- Else ($E$): Choose between Latency ($L$) or Consistency ($C$).
- Examples:
- PC/EC (e.g., Spanner, CockroachDB): Prefers consistency both in partition and normal operation.
- PA/EL (e.g., DynamoDB, Cassandra with weak consistency): Prefers availability during partition and low latency during normal operation.
03 — The Raft Consensus Algorithm
Raft decomposes distributed consensus into three understandable, orthogonal sub-problems:
[Follower] ──(Election Timeout)──> [Candidate] ──(Votes from Majority)──> [Leader]
▲ │
└──────────────────────(Discovers Higher Term)───────────────────────────┘
- Leader Election: When heartbeat timers expire, a follower transitions to candidate, increments the current term, votes for itself, and broadcasts RequestVote RPCs. A candidate winning votes from a majority of nodes becomes Leader.
- Log Replication: The leader receives client commands, appends them to its local log, and broadcasts AppendEntries RPCs. Once an entry is replicated across a majority of nodes, the leader commits it and applies it to its state machine.
- Safety: Raft guarantees that if a leader commits a log entry at a given index and term, no other server will ever apply a different log entry for that index.