9.8 Technology deep dives
φ = −log₁₀(P(heartbeat arrives later than the current elapsed time)).
8.1 Failure detectors (timeouts, heartbeats, Phi Accrual)
Problem it solves. Decide whether a node is dead, so a load balancer can remove it or a follower can be promoted — without any reliable way to observe the difference between "dead," "slow," and "unreachable."
Why wasn't an explicit signal enough? §2.4: RST/FIN, crash scripts, switch queries, and ICMP all exist and none of them can be relied upon. In general you get no response at all.
Why wasn't a fixed timeout enough? §2.5: there is no correct value. Too short → false positives → load transfer → cascading failure. Too long → slow recovery. And the delay distribution is not stationary — it depends on load, on noisy neighbours, on whether a switch is being upgraded.
How Phi Accrual works internally. Instead of a boolean "up/down," it maintains a sliding window of recent heartbeat inter-arrival times, fits a distribution (normal in the classic formulation), and computes
φ = −log₁₀(P(heartbeat arrives later than the current elapsed time)).
φ rises continuously as the silence lengthens. The application picks a threshold: φ ≥ 8 means "the probability we are wrong is about 10⁻⁸." The value is that the detector adapts automatically to a network whose latency distribution changes — exactly the §2.5 recommendation.
Deployment. Heartbeats between all peers (Cassandra's gossip) or from a leader; thresholds per-role — be more conservative about leader failover than about removing a node from a read pool, because the cost of a false positive differs by an order of magnitude.
Monitoring. False-positive rate (nodes marked dead that recover within seconds — the single most useful signal that your threshold is too aggressive); time-to-detection distribution; heartbeat inter-arrival p99; flapping count per node; the correlation between detection events and load spikes (if they correlate, you are detecting overload, not death).
Scaling. All-to-all heartbeating is O(n²); above a few hundred nodes you need gossip or a hierarchy.
What actually breaks.
- Detecting overload as death, then transferring load away, then detecting more overload — §2.5's cascading failure, and the reason Ch 7 warned against automatic rebalancing plus automatic failure detection.
- A GC pause exceeding the timeout, so a perfectly healthy leader is demoted mid-write.
- Asymmetric faults (§2.3): the node is alive and receiving, but its outbound path is broken, so it is declared dead while it keeps believing it is the leader — the §5.1 parable, and precisely the situation that requires fencing.
- Health checks that only prove the process is running, not that it can do work — so a fail-slow / limping node stays in rotation forever, which §6.2 notes is harder to handle than a clean crash.
- Flapping when the threshold sits right at the network's p99.
8.2 NTP, PTP, and clock monitoring
Problem it solves. Give every machine roughly the same notion of wall-clock time, so timestamps in logs, expiries, and cross-machine reasoning are comparable.
Why wasn't the hardware clock enough? §3.2 ①: up to 200 ppm drift, temperature-dependent — 17 seconds per day if never resynced.
Why isn't NTP enough for ordering? §3.4: NTP's accuracy is bounded by network round-trip time, so you cannot make clock error smaller than network delay, which is exactly what correct ordering would require.
How it works internally. NTP samples several servers, measures offset and round-trip delay for each, discards outliers (the "weak lying" defence of §5.3), and either slews (gradually adjusts the rate, ≤0.05%) or steps (jumps) the clock. PTP pushes accuracy to sub-microsecond by using hardware timestamping in the NIC and transparent clocks in switches that account for their own queueing delay — which is why PTP needs switch support and NTP doesn't.
Deployment. Stratum-1 sources (GPS/atomic) in each datacenter; at least 4 upstream servers so outlier rejection works; chrony over ntpd on modern Linux (faster convergence, better on VMs); leap-second smearing configured consistently across the fleet — a mixed fleet where half smear and half step is worse than either.
Monitoring — and this is the section people skip.
- Clock offset per node vs the reference — with an alert threshold well below any application-level assumption
- Whether the NTP daemon is actually running and synchronized (
chronyc tracking: stratum, root dispersion, last offset). §3.2 ③: a node firewalled off from NTP looks perfectly healthy while drifting - Step (jump) events — every one of these is a potential correctness incident
- Root dispersion / estimated error — this is your confidence interval, and §3.5 says most systems never look at it
- On VMs: steal time, since §3.2 ⑦ means the guest's own accuracy estimate can be wrong
Backup. Not applicable, but: have a plan for GPS jamming (§3.2), which is a real and locality-dependent risk.
What actually breaks.
- Silent drift after a firewall change — the canonical §3.3 failure: nothing errors, and data quietly disappears via LWW.
- A leap second hanging every JVM in the fleet simultaneously (the Ch 2 example) — a correlated software fault.
- VM live migration making the clock jump forward mid-transaction.
- Half the fleet smearing and half stepping during a leap second, so nodes disagree by a full second for a day.
- Using
System.currentTimeMillis()to measure a duration, so a backward step produces a negative elapsed time — and code that divides by it. - Trusting microsecond digits when root dispersion is 40 ms (§3.5).
8.3 Distributed lock services (ZooKeeper, etcd, Chubby) and fencing
Problem it solves. Make exactly one node the leader / lock holder, and — crucially — make it safe when that guarantee is inevitably violated.
Why wasn't a lock in a database enough? A lock with no lease expires never (the holder crashes and blocks forever); a lock with a lease can be held by two nodes at once (§5.2 cases ① and ②).
Why isn't STONITH enough? §5.2: it doesn't stop a delayed packet, all nodes can shoot each other, and it may act too late.
How it works internally. A consensus-replicated state machine (Ch 10). Locks are ephemeral nodes / leases tied to a session with a heartbeat; when the session lapses, the lock is released automatically. Every state change carries a monotonically increasing version — ZooKeeper's zxid/cversion, etcd's revision — which is what makes it a fencing-token generator, satisfying §6.3's uniqueness and monotonic sequence safety properties.
The critical design rule: the lock service alone is not sufficient. The RESOURCE must check the token. A lock service without token-checking on the storage side gives you the illusion of mutual exclusion and none of the substance.
Deployment. 3 or 5 nodes (odd, for majority quorums), spread across availability zones but NOT across regions — cross-region consensus latency makes every lock acquisition painful. Session timeout tuned above your worst realistic GC pause.
Monitoring. Session expiry rate (each one is a potential zombie); leader elections per hour (frequent elections = your timeouts fight your GC pauses); request latency p99; fsync latency on the consensus log — this bounds everything; watch-count and znode-count (ZooKeeper falls over on unbounded watches); quorum health and whether any member is lagging.
Backup. etcd snapshots / ZooKeeper transaction-log + snapshot backups. Losing the coordination store loses the shard map, the leader identity, and every lock — and Ch 7 §4 notes it's the authority for routing too.
What actually breaks.
- A lock service used without fencing tokens. The system looks correct in testing and corrupts data in production, exactly as HBase did.
- A session timeout shorter than a GC pause → the leader is demoted while alive → split brain until fencing catches it.
- A client that acquires a lease and then does a long operation without re-checking, so §5.2 case ① is guaranteed rather than merely possible.
fsynclatency spikes on the consensus log stalling all lock operations cluster-wide.- Cross-region deployment turning every lock acquisition into a 150 ms round trip.
- Treating the lock service as available: it is CP, not AP — during a partition the minority side cannot acquire locks, by design.
8.4 Jepsen, TLA+, and deterministic simulation (Antithesis, FoundationDB)
Problem it solves. Establish that an implementation actually satisfies its claimed safety properties under the faults of §§2–4 — which normal testing cannot do, because the bugs live in orderings you never happened to produce.
Why wasn't unit/integration testing enough? §7: the state space is enormous, and concurrency bugs only manifest when you get unlucky with the timing (Ch 8 §3).
Why isn't one technique sufficient?
| Technique | Finds | Misses |
|---|---|---|
| TLA+ / model checking | Design-level bugs, protocol ambiguities (the viewstamped-replication data-loss finding) | Anything where the implementation diverges from the spec |
| Jepsen / fault injection | Real bugs in the real system under real faults | Not replayable; coarse control; only explores orderings it happens to hit |
| DST | Real bugs in real code, replayable, systematic | Requires controlling every source of nondeterminism — an architectural commitment |
How Jepsen works internally. Runs the real cluster; a generator issues concurrent client operations while a nemesis injects faults (partition, clock skew, process kill, pause); every operation is recorded with invoke/complete times into a history; a checker (Knossos/Elle) then searches for a linearization or a serializable schedule consistent with that history — if none exists, it produces a concrete counterexample. Elle goes further and infers transaction dependency cycles, naming the exact Ch 8 anomaly observed.
Monitoring/practice. Run in CI on a schedule, not once before launch; keep the failing histories — they're the regression suite; treat a Jepsen finding as a design review trigger, not just a patch.
What actually breaks.
- Spec/implementation drift — a beautifully verified TLA+ model of code that no longer matches it (§7.1).
- DST that misses a nondeterminism source — hash iteration order, allocation failure (§7.4) — so a "deterministic" replay isn't.
- Testing only clean crashes and never latency, partial partitions, or clock skew — which are the faults that actually break systems.
- Running chaos/fault injection and not fixing what it finds, converting the practice into theatre.