Learn Labs
10. Consistency and Consensus

10.5 Technology deep dives


5.1 Raft (etcd, CockroachDB, TiKV, Consul, RabbitMQ quorum queues)

Problem it solves. Replicate a log across a set of nodes with automatic, safe leader failover — so that no acknowledged write is ever lost and split brain is impossible.

Why wasn't single-leader replication enough? It has no safe automatic failover: promoting a lagging follower loses committed writes; two nodes can both believe they're leader (Ch 6 §1.4).

Why wasn't Paxos enough? Raft was explicitly designed for understandability — Paxos's Multi-Paxos extension is what people actually deploy, and it's notoriously underspecified in the literature. Raft prescribes leader election, log replication, and membership change concretely.

How it works internally. Three roles (follower / candidate / leader). Terms are the epoch numbers: monotonically increasing, at most one leader per term. Election timeout randomized (typically 150–300 ms) to avoid split votes. Election restriction: a candidate only wins if its log is at least as up-to-date as the voter's — this is Raft's version of §3.5's "new leader must be up to date." A leader appends to its log, replicates via AppendEntries (which doubles as the heartbeat), and commits an entry once a majority have it — and only entries from its own current term directly, which is the subtle rule that prevents a committed entry from being overwritten. PreVote was added to fix the §3.6 leadership-bouncing pathology.

Deployment. 3 or 5 members (odd, so a majority exists), spread across AZs but rarely across regions — every write costs one quorum round trip, so cross-region raises write latency to the inter-region RTT. Learners/non-voting members for adding capacity or a new region without changing the quorum.

Monitoring.

  • Leader elections per hour — should be ~0. Nonzero and rising means your election timeout is fighting your GC pauses or network jitter (§3.6 cost #4)
  • Commit latency p99 and fsync latency on the WAL — the latter bounds everything
  • Follower lag / matchIndex gap per member — a lagging follower is a member that can't become leader safely
  • Proposal failure rate; quorum-loss events
  • etcd specifically: DB size vs quota (a full etcd goes read-only — a classic Kubernetes outage), compaction and defrag status, watch count

Backup. Periodic snapshots + the log. Losing a majority of members loses the cluster, so snapshots must live off-cluster.

What actually breaks.

  • Election storms. Timeouts too aggressive relative to real network/GC behavior; the cluster spends more time electing than working — verbatim §3.6.
  • fsync latency spikes on a slow disk stalling every write cluster-wide, because commit requires a durable quorum.
  • etcd exceeding its storage quota and flipping to read-only, taking Kubernetes with it.
  • Cross-region deployment making every write pay an inter-continental RTT.
  • Losing quorum (2 of 3 members) — the cluster is safe but unavailable, and recovery requires a deliberate, dangerous force-new-cluster operation.
  • Reading from a follower and assuming linearizability — §1.5's warning: consensus doesn't make all operations linearizable.

5.2 ZooKeeper (Zab) and the coordination-service pattern

Problem it solves. Give other distributed systems the primitives they need — leases, fencing tokens, failure detection, change notification — without each of them implementing consensus.

Why not embed consensus in every system? §4: with thousands of shards, running consensus over thousands of nodes is terribly inefficient. Outsource it to 3–5 nodes.

How it works internally. Zab (a Paxos-family protocol with a strong ordering guarantee). A hierarchical namespace of znodes; ephemeral znodes tied to a session (the failure detector); sequential znodes producing monotonic counters (the fencing tokens); watches as one-shot change notifications. zxid is a 64-bit value: high 32 bits = epoch, low 32 bits = counter within the epoch — so it is simultaneously a fencing token and a leadership generation marker.

Deployment. Ensemble of 3 or 5; observers for read scaling and remote regions (§4). Session timeout tuned above worst-case GC pause. Apache Curator for recipes (leader election, distributed locks, barriers) — because the raw API is easy to misuse.

Monitoring. Outstanding requests; session expirations per hour (each is a potential zombie — Ch 9 §5.2); watch count (unbounded watches are a known way to melt an ensemble); znode count and data size; fsync time; leader election count; follower sync latency.

What actually breaks.

  • Using ZooKeeper reads as linearizable. §1.4's footnote: writes are linearizable, reads may be stale. You must issue a sync before a read if you need recency.
  • Watch semantics misunderstood — watches are one-shot and may miss intermediate states; code that assumes it sees every change is wrong.
  • Session timeout shorter than a GC pause → the ephemeral node vanishes → the lease is reassigned → zombie.
  • Storing too much or too fast-changing data — §4's constraint. This is not a database.
  • Watch explosion from a client registering a watch per key across a large keyspace.
  • A herd effect when a leader's ephemeral node disappears and every candidate wakes at once (Curator's recipes exist to avoid exactly this).

5.3 Spanner / TrueTime — linearizability without a coordination round trip

Problem it solves. Globally distributed strict serializability where a single-node timestamp oracle (§2.4) would force every transaction in every region through one place.

Why not a timestamp oracle? §2.4's limits: can't shard it, can't distribute it across regions — in a geo-distributed database, all ID requests would go to a node in a single region.

Why not HLCs? §2.3: they give causal ordering, not linearizability — the privacy-leak example is exactly the failure.

How it works internally. TrueTime returns [earliest, latest] from GPS + atomic clocks per datacenter, keeping ε ≈ 7 ms (Ch 9 §3.6). A read/write transaction picks a commit timestamp and then waits out ε before releasing locks ("commit wait"), guaranteeing that any later transaction gets a strictly greater timestamp without any communication. Reads at a timestamp are lock-free and served from any sufficiently up-to-date replica.

Monitoring. TrueTime ε (if it grows, commit-wait grows and every write slows — this is a hardware health metric with a direct latency consequence); commit-wait duration; per-Paxos-group leader locality; transaction retry/abort rate; participant count per transaction (2PC over Paxos groups: more groups = more expensive).

What actually breaks. ε widening after a GPS or time-master failure, silently taxing every write; cross-region transactions costing consensus round trips plus commit wait; and the general trap of assuming a globally distributed transaction is as cheap as a local one.


5.4 Distributed ID generation (Snowflake, UUIDv7, ULID, timestamp oracles)

Problem it solves. Unique primary keys, generated without a bottleneck, ideally sortable by creation time.

Why not autoincrement? §2: SPOF, cross-region latency, throughput bottleneck.

Why not UUIDv4? §2.1: 128 bits and random order — comparing two tells you nothing about which is newer, which also destroys B-tree insert locality (Ch 4 §9.2: random inserts into a clustered index cause page splits and poor fill factor).

How Snowflake-style works internally. [41 bits ms timestamp][10 bits machine ID][12 bits sequence] = 64 bits, k-sortable, ~4,096 IDs/ms/node. UUIDv7 is the standardized version of the same idea: 48-bit Unix ms timestamp in the high bits, then random. ULID is the same shape with Crockford base32 text encoding.

Monitoring. Sequence exhaustion (>4,096/ms on one node → the generator must wait or error); clock rollback events — the failure that forces a Snowflake node to refuse to issue IDs; machine-ID collisions (the nastiest failure, and easy to cause with autoscaling that reuses IDs).

What actually breaks.

  • Clock moving backward — Snowflake's only safe response is to stop issuing IDs until the clock catches up, which is a hard outage caused by NTP.
  • Duplicate machine IDs after an autoscaler recycles an instance ordinal → duplicate primary keys, discovered much later.
  • Assuming timestamp-prefixed IDs are causally ordered. They are approximately time-ordered and not linearizable (§2.1) — which is precisely the privacy-leak trap.
  • Hot shard from monotonic keys (Ch 7 §3.1): the very sortability that helps B-trees sends all writes to one shard.

On this page