Chapter 08 · Membership, Failure Detection, Gossip, and Cluster Coordination
Heartbeats, Phi Accrual and Other Failure Detectors: Suspicion Is Not Proof
Treat heartbeat timeouts and φ-style outputs as suspicion under uncertainty, then show how adaptive evidence differs from a brittle binary timeout.
Learning outcomes
AtlasMart's node-c stops answering heartbeats for
1.6 seconds during a storage pause, then recovers. Another node
sees packet loss on only one network path. If the cluster
equates “missed heartbeat” with “dead,” it can trigger avoidable
replica movement or leadership changes while the process is
still alive.
Distinguish heartbeat evidence, timeout policy, suspicion, declaration, and recovery.
Explain false positives and the latency-versus-certainty tradeoff of failure detection in asynchronous networks.
Compute and interpret an accrual-style suspicion score without treating it as proof of failure.
Reason about asymmetric reachability: A can suspect B while C still communicates with B.
Connect detector output to safe actions such as routing avoidance, quorum decisions, and fencing rather than destructive automation.
A node crash, a congested link, packet loss, scheduler pause, long garbage-collection pause, overloaded disk, and firewall mistake can all look like “no reply before my deadline.” From local timing alone, a detector cannot identify the cause with certainty.
1. Heartbeats turn silence into evidence, not knowledge
A heartbeat is a periodic liveness message or probe/acknowledgement exchange. A fixed detector might declare a peer down after k missed heartbeats or after a fixed timeout. This is easy to reason about but brittle when network and scheduler delay vary.
A false positive occurs when a healthy/recoverable member is classified as failed. Increasing the timeout lowers false positives but increases detection latency. Decreasing it reacts faster but can mistake ordinary tail latency for failure. This is a policy tradeoff, not a magic constant.
| Observation at node-a | Possible reality | Safe interpretation |
|---|---|---|
| No heartbeat for 1.5 s | node-b crashed | suspect |
| No heartbeat for 1.5 s | a→b path congested | suspect |
| No heartbeat for 1.5 s | node-b paused on disk/GC/scheduler | suspect |
| c still receives b | asymmetric reachability | do not claim globally dead from a alone |
2. Accrual detectors separate measurement from policy
An accrual failure detector exposes a continuously increasing suspicion value rather than a single Boolean answer. The application then chooses a threshold appropriate to its action. A latency-sensitive request router may avoid a node at a lower suspicion level than the threshold used for a disruptive membership removal.
The φ (phi) accrual detector expresses suspicion roughly as the negative base-10 logarithm of the probability that an arrival would be later than the current silence, under a learned arrival-time model. A larger φ means the observed silence is increasingly surprising relative to recent history. The exact estimator, sampling window, minimum variance, pause handling, and threshold are implementation details.
A φ value is meaningful only relative to the detector implementation, its history, workload, network, and policy. A threshold copied from another product or environment is not a portable reliability guarantee.
3. Asymmetric reachability changes the operational picture
t=0.0 a,b,c exchange probes normally
t=1.0 path a -> b begins dropping packets
t=1.1 c -> b succeeds; b -> c succeeds
t=1.5 a's suspicion(b) rises
t=1.6 client near a times out to b
t=1.7 c still serves traffic through b
conclusion:
"a suspects b" != "b is globally dead"
This is partial failure from Chapter 02 expressed through detector state. The detector output should feed a protocol that can tolerate disagreement—for example, choose another coordinator, wait for quorum evidence, or require an ownership/term transition before accepting a replacement owner.
4. Wrong approach: trigger destructive repair on the first timeout
Automatically replacing or decommissioning a member after one application timeout can make a transient problem worse. The system may stream large replicas, increase disk/network pressure, invalidate caches, and cause a second wave of timeouts—the classic self-inflicted failure cascade.
A safer ladder is action-specific: retry/read from another replica; mark the path degraded; accumulate detector evidence; corroborate with peers; preserve the old owner's fencing/term; and only perform topology-changing repair under documented membership rules. If the old owner later returns, it must not regain authority merely because it is reachable again.
5. AtlasMart lab — fixed timeout versus suspicion score
The simulator uses a normal-distribution approximation only to make the idea visible; it is not an implementation of Cassandra or any production φ detector.
from math import erf, log10, sqrt
from statistics import mean, pstdev
# Heartbeat arrival intervals in seconds observed before a delay spike.
history = [0.92, 1.08, 1.02, 0.97, 1.11, 0.95, 1.04, 0.99]
mu = mean(history)
sigma = max(pstdev(history), 0.05)
def normal_cdf(x, mu, sigma):
return 0.5 * (1.0 + erf((x - mu) / (sigma * sqrt(2.0))))
def phi(elapsed):
# Pedagogical normal-distribution approximation, not a product algorithm.
survival = max(1.0 - normal_cdf(elapsed, mu, sigma), 1e-12)
return -log10(survival)
print(f"history mean={mu:.3f}s sigma={sigma:.3f}s")
for elapsed in [1.10, 1.25, 1.40, 1.55, 1.80]:
fixed_timeout_dead = elapsed > 1.30
suspicion = phi(elapsed)
print(f"elapsed={elapsed:.2f}s fixed_dead={fixed_timeout_dead} phi={suspicion:.2f}")
print("\nInterpretation:")
print("- fixed 1.30s timeout turns lateness into a binary failure claim")
print("- phi is a suspicion score; policy chooses a threshold")
print("- even a high phi cannot prove whether node, path, scheduler, or disk is the cause")
Expected output
history mean=1.010s sigma=0.061s
elapsed=1.10s fixed_dead=False phi=1.16
elapsed=1.25s fixed_dead=False phi=4.40
elapsed=1.40s fixed_dead=True phi=10.14
elapsed=1.55s fixed_dead=True phi=12.00
elapsed=1.80s fixed_dead=True phi=12.00
Interpretation:
- fixed 1.30s timeout turns lateness into a binary failure claim
- phi is a suspicion score; policy chooses a threshold
- even a high phi cannot prove whether node, path, scheduler, or disk is the cause
The fixed detector crosses its 1.30-second line abruptly. The accrual score grows continuously. Neither output identifies the root cause, and neither by itself authorizes a stale owner or destructive membership change.
Check your understanding
- Why can no heartbeat be ambiguous?
- What does an accrual detector add?
- Why not copy one φ threshold everywhere?
- What evidence can reduce an asymmetric false positive?
- What must still protect data if a stale node resumes?
Review the answers
1. Because crash, network loss, overload, pauses, and asymmetric paths can produce the same local observation.
2. A graded suspicion value separated from the application policy that chooses when/how to act.
3. Its false-positive/detection behavior depends on estimator, history, environment, and the consequence of the action.
4. Peer observations, alternate paths, quorum/ownership state, and recovery of the suspected node.
5. Ownership epochs/terms or fencing so reachability does not automatically restore write authority.
6. Production judgment
Monitor detector suspicion, heartbeat interval distributions, network tail latency, false-positive rate, flap count, coordinator retries, membership churn, and the expensive work triggered by suspicion. Test detector behavior under injected delay and process pauses in an isolated simulator or test cluster—not by globally breaking the learner's network.
The core guarantee is modest but important: a detector helps the system make progress by producing evidence under uncertainty. It does not turn an asynchronous network into a failure oracle. The next lesson examines how membership observations spread through the cluster after one node learns something new.
Authoritative references
- The φ accrual failure detector — Hayashibara et al.; graded suspicion abstraction for failure detection
- SWIM membership protocol — uses suspicion before declaration to reduce false positives
- Cassandra failure-detector tool — current product example for observing detector state; not a universal algorithm definition