Part 5 · 2 chapters · ~12 min
Distributed Theory
CAP stated formally (Gilbert and Lynch, 2002) and the cartoon version debunked, PACELC, consistency models as a hierarchy, the FLP impossibility result, consensus (Paxos, Raft) and what it solves, quorum maths (R + W > N, sloppy quorums, hinted handoff), Lamport, vector and hybrid logical clocks, CRDTs (state-based and operation-based, semilattices), and two-phase commit and why it blocks.
10
CAP, PACELC, FLP
| result | precise statement | engineering meaning |
|---|---|---|
| CAP (Gilbert-Lynch 2002) | no system provides linearisability and availability (every non-failed node responds) under network partitions | during a partition, choose: refuse some requests or allow stale or divergent answers. It is not "pick two of three" in normal operation. |
| PACELC (Abadi 2012) | if Partition, trade A against C; Else, trade Latency against Consistency | the everyday trade-off is latency, not partitions: synchronous replication costs round trips |
| FLP (1985) | no deterministic protocol solves consensus in an asynchronous system if even one process may crash | real systems use timeouts (partial synchrony) and randomness: Raft is safe always, live only when timing behaves |
CONSISTENCY MODELS AS A HIERARCHY
stronger models allow fewer surprising histories
swipe the figure sideways, or tap expand for full screen
1/4
linearisable
The strongest single-object model: once a write completes, every later read sees it. Required for locks, leader election and balance checks; costs coordination on every operation.
real-time orderetcd, Spanner, single leader
11
Quorums, clocks, CRDTs, 2PC
code
quorums: N replicas, write to W, read from R. R + W > N ⇒ every read quorum overlaps the latest write quorum
N = 3, W = 2, R = 2 → 2 + 2 > 3 ✓ N = 3, W = 1, R = 1 → stale reads possible
sloppy quorums + hinted handoff (Dynamo): accept writes on stand-in nodes during failures → no overlap guarantee
clocks: Lamport: on send/receive, t = max(local, received) + 1 → consistent with causality, cannot detect concurrency
vector: one counter per node; a ≤ b component-wise ⇔ a happened before b; incomparable ⇔ concurrent
HLC: physical time + logical counter: close to wall clock, still causal (CockroachDB, MongoDB)
CRDTs: state merges must form a join-semilattice: commutative, associative, idempotent
G-Counter: merge = element-wise max · PN-Counter = two G-Counters · OR-Set: adds win over concurrent removes
2PC: coordinator: PREPARE → all YES → COMMIT. A participant that voted YES must wait for the decision;
if the coordinator dies then, it is blocked holding locks. 3PC removes blocking only under synchronous assumptions;
practical fixes replicate the coordinator with consensus (Spanner) or avoid 2PC with sagas (Workflows course).