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

resultprecise statementengineering meaning
CAP (Gilbert-Lynch 2002)no system provides linearisability and availability (every non-failed node responds) under network partitionsduring 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 Consistencythe 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 crashreal 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
linearisableeach operation appears instant, in real-time ordersequentialone global order consistent with each client's ordercausalcausally related operations seen in order by allsession guaranteesread your writes, monotonic readseventualreplicas converge if writes stop
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).