Every distributed system that reacts to a dead node, whether it reroutes traffic, rebalances partitions or elects a new leader, depends on a component that answers one question: is that process still there? That component is a failure detector, and the uncomfortable truth behind it is that in a real network the question has no reliable answer. A process that has crashed and a process whose replies are merely slow look identical from outside, because silence is the only evidence either one produces.
This article treats failure detection as engineering rather than a timeout setting: the theory of what a detector can promise, the metrics to judge one by, runnable phi accrual code in two models, SWIM-style probing, and the failure modes behind real outages. For the phi pipeline alone, the companion piece on phi accrual as a signal is shorter; this page goes wider and into code.
Why a detector can only ever suspect
In an asynchronous system there is no upper bound on message delay or on how long a process may pause. Fischer, Lynch and Paterson showed in 1985 that under those assumptions no deterministic protocol can guarantee consensus if even one process may crash. The intuition is exactly the detection problem: a protocol waiting on a slow process can never safely decide that it has crashed, and a protocol that stops waiting may exclude a process that is alive and about to act.
Chandra and Toueg turned this into a design tool in 1996. They modelled a failure detector as a local oracle that outputs a list of suspected processes, and they classified detectors by two properties. Completeness asks whether every crashed process is eventually suspected; strong completeness means by every correct process. Accuracy asks how often correct processes are wrongly suspected. Strong accuracy means never; weak accuracy means at least one correct process is never suspected; the eventual variants only require these to hold after some unknown time.
Combining strong completeness with each accuracy level gives the classes P (perfect), S (strong), ◇P (eventually perfect) and ◇S (eventually strong); weak completeness gives Q, W, ◇Q and ◇W. Chandra, Hadzilacos and Toueg later proved that ◇W, which is equivalent to the leader oracle Ω, is the weakest detector with which consensus can be solved when a majority of processes is correct. That is why Paxos and Raft are built so that a detector may be wrong for a while without breaking safety: they only need it to be right eventually, and they protect correctness with quorums rather than with suspicion.
The practical lesson is blunt. No timeout, however clever, turns suspicion into knowledge. A detector is a performance component that decides how quickly you react and how often you react wrongly. Anything that must be correct, such as a single writer or a single leader, needs a separate mechanism like a time-bounded lease or fencing tokens.
Measuring a detector: the three numbers that matter
Chen, Toueg and Aguilera proposed quality-of-service metrics in 2002 that are far more useful than talking about a timeout value. They separate speed from mistakes, which lets you compare detectors that are configured completely differently.
| Metric | Meaning | What moves it |
|---|---|---|
| Detection time TD | Time from a real crash until the detector permanently suspects the process | Heartbeat interval plus the threshold or timeout |
| Mistake recurrence time TMR | Average time between two wrong suspicions of a live process | Network jitter, pauses, threshold height |
| Mistake duration TM | How long a wrong suspicion lasts before it is corrected | Heartbeat interval, refutation speed |
Tuning is always a trade between TD and TMR. Halving the heartbeat interval improves both but doubles background traffic, which in an all-to-all scheme grows with the square of cluster size. Raising the threshold improves TMR and hurts TD. A tuning change that only reports faster detection is hiding its cost.
The detector family
Four designs cover almost everything you will meet.
| Design | Decision rule | Strength | Weakness |
|---|---|---|---|
| Fixed timeout | Suspect if silence exceeds D | Trivial, predictable TD | One D for every link; wrong in both directions at once |
| Adaptive estimate (Chen) | Suspect if silence exceeds the estimated next arrival plus a safety margin | Tracks the mean delay of each link | Margin is still a constant; ignores variance |
| Phi accrual (Hayashibara et al.) | Output phi = -log10 P(heartbeat arrives later than now); consumers compare with their own threshold | Adapts to variance; one stream, many policies | Assumes a distribution; bad under bursty pauses |
| SWIM probing (Das, Gupta, Motivala) | Probe one random member per period; ask k helpers to probe indirectly before suspecting | Constant per-node load; filters out a bad local path | Detection is probabilistic; needs suspicion and refutation logic |
Phi accrual separates estimation from decision: it fits a distribution to recent inter-arrival gaps and reports how surprising the current silence is on a log scale, where each unit is another factor of ten. A router can act at a low phi; a failover controller waits for a high one.
Building phi accrual, both models
The code below is a complete, dependency-free detector for a single peer. It supports the normal model used by Akka, including Akka's two guard rails: a floor on the standard deviation, and an acceptable pause that is added to the mean. It also supports the exponential model used by Cassandra, where phi is simply the elapsed time divided by the mean gap, multiplied by log10(e). Cassandra's default phi_convict_threshold is 8; Akka's default threshold is 8.0 and its cluster failure detector adds an acceptable heartbeat pause of 3 seconds by default.
import math
from collections import deque
class PhiDetector:
"""Phi accrual over one monitored peer. Times are monotonic milliseconds."""
def __init__(self, model="normal", window=1000, min_std_ms=100.0,
acceptable_pause_ms=0.0, first_interval_ms=1000.0):
self.model = model
self.gaps = deque(maxlen=window)
self.min_std = min_std_ms
self.pause = acceptable_pause_ms
self.last = None
# Seed with a spread so early variance is not near zero.
self.gaps.extend([first_interval_ms - first_interval_ms / 4,
first_interval_ms + first_interval_ms / 4])
def heartbeat(self, now_ms):
if self.last is not None:
self.gaps.append(now_ms - self.last)
self.last = now_ms
def phi(self, now_ms):
if self.last is None:
return 0.0
elapsed = now_ms - self.last
mean = sum(self.gaps) / len(self.gaps)
if self.model == "exponential": # Cassandra-style: phi grows linearly
return elapsed / mean * math.log10(math.e)
var = sum((g - mean) ** 2 for g in self.gaps) / len(self.gaps)
std = max(math.sqrt(var), self.min_std) # floor: stop a quiet LAN over-reacting
return _phi_normal(elapsed, mean + self.pause, std)
def _phi_normal(t, mean, std):
"""Akka's logistic tail approximation; both of its branches equal log10(1 + exp(z))."""
y = (t - mean) / std
z = y * (1.5976 + 0.070566 * y * y)
if z > 0:
return z / math.log(10) + math.log10(1 + math.exp(-z))
return math.log1p(math.exp(z)) / math.log(10)Three details matter more than the formula. Use a monotonic clock, or an NTP step will convict the whole cluster. Seed the window, because two samples give a variance near zero. And compute the tail carefully: Akka's formula written literally as -log10(e / (1 + e)) underflows to zero once the silence is far past the mean, which returns infinity on the JVM and raises an exception in Python. Both branches equal log10(1 + ez), and the version above uses that identity to stay finite in both directions.
Worked example: what the numbers look like
Suppose peer A sends a heartbeat every second and the window has settled at a mean gap of 1,000 ms with a standard deviation of 100 ms. Using the normal model with no acceptable pause, the detector code produces these values:
| Silence since last heartbeat | phi (normal, sd 100 ms) |
|---|---|
| 1,000 ms | 0.30 |
| 1,200 ms | 1.64 |
| 1,400 ms | 4.74 |
| 1,500 ms | 7.30 |
| 1,600 ms | 10.78 |
Phi crosses 8 at about 1,523 ms. That is aggressive: a single stop-the-world pause of 600 ms in either process would produce a false suspicion. This is exactly why Akka adds an acceptable pause. With 3,000 ms added to the mean, the same threshold is crossed at about 4,523 ms. Widening the observed standard deviation to 500 ms, as a noisy cross-region link would, moves the crossing to about 3,613 ms without any configuration change, which is the adaptivity phi is known for.
The exponential model behaves differently. Because phi equals elapsed divided by mean times 0.434, reaching phi 8 needs a silence of 8 times ln 10, about 18.4 times the mean gap. With a one-second gossip interval that is roughly 18 seconds, and the curve is linear, so it ignores variance entirely. The two models with the same threshold of 8 therefore mean very different detection times. Never copy a threshold between systems without recomputing what it implies.
SWIM: probing instead of listening
All-to-all heartbeats cost O(n2) messages per interval, which stops scaling at a few hundred nodes. SWIM keeps per-node load constant. In each protocol period, a member pings one target chosen in round-robin order over a shuffled list. If no acknowledgement arrives, it asks k other members to ping the target on its behalf. Only if all of those fail is the target marked suspect, and the suspicion is spread by piggybacking on the probe messages themselves.
import random
def probe_round(me, members, send_ping, send_ping_req, k=3, timeout_ms=500):
"""One SWIM protocol period: direct probe, then k indirect probes, then suspect."""
target = members.next_in_shuffled_order(exclude=me) # round-robin over a shuffle
if send_ping(target, timeout_ms):
return ("alive", target)
pool = members.alive(exclude={me, target})
helpers = random.sample(pool, min(k, len(pool)))
acks = [send_ping_req(h, target, timeout_ms) for h in helpers] # in parallel in practice
if any(acks):
return ("alive", target) # our path was bad, not the target
members.mark_suspect(target, incarnation=members.incarnation(target))
return ("suspect", target) # gossiped; confirmed dead only after a suspicion timeoutIndirect probing is the key accuracy improvement: it separates a dead target from a broken path between the prober and the target. The suspicion state is the second: a suspected member has a timeout during which it can refute the claim by gossiping a higher incarnation number. HashiCorp's Lifeguard extensions, used in memberlist, add local health awareness, so a member that is itself slow, for example because it is missing its own acknowledgements, stretches its timeouts instead of accusing healthy peers. The article on gray failure explains why that self-doubt matters.
Failure modes that cause real incidents
- Observer pauses. A long garbage collection or VM steal on the observer makes every peer look late at once. Detect a local pause by comparing scheduler ticks with the monotonic clock, then discard the affected samples instead of convicting everyone.
- Heartbeats on a busy thread. If heartbeats are sent from an executor that also serves requests, load delays them, and the detector reports overload as death. Send and receive heartbeats on a dedicated thread or channel with priority.
- Gray failure. The heartbeat path is healthy while the data path is broken, such as a full disk or a stuck request queue. The detector says alive while clients time out. Feed request success rates into health decisions, not only heartbeats.
- Asymmetric partitions and suspicion storms. A can hear B but B cannot hear A, or one switch reboot makes many nodes suspect many others. Membership flaps and rebalancing moves data when the cluster is weakest. Rate-limit reactions and require several observers to agree.
- Tiny variance. On a quiet LAN the measured standard deviation can fall to a few milliseconds, turning any hiccup into a high phi. Always floor it.
- Acting on suspicion as if it were proof. Promoting a new primary because phi crossed a threshold, without fencing the old one, is how split brain happens.
Operating and tuning a detector
Start by deciding what each consumer does with a suspicion and what a wrong one costs. Export phi, or the silence duration, per peer as a metric and look at the per-minute maximum during normal operation. Your threshold should sit comfortably above the tail of that distribution, not above its median.
In Cassandra, operators on noisy virtualised networks often raise phi_convict_threshold above the default; because the model is exponential, each additional unit adds about 2.3 mean gaps of silence to detection time, which you can compute before you change it. In Akka, prefer increasing acceptable-heartbeat-pause to handle known pauses, such as garbage collection, and leave the threshold alone. For consensus systems, keep the election timeout far above typical heartbeat jitter; the Raft consensus article shows how randomised timeouts avoid split votes.
Finally, test deliberately: pause a process with SIGSTOP, add latency with a traffic shaper, drop packets in one direction, and keep the measured TD and false suspicions as a baseline.
What to do next
- Write down, per consumer, what happens on a suspicion and what a false positive costs.
- Instrument heartbeat inter-arrival gaps and per-peer phi or silence, and look at the p99.9 of the per-minute maximum over a normal week.
- Compute the detection time your current threshold implies with the code above, for your real mean and standard deviation.
- Move heartbeats onto a dedicated thread and use a monotonic clock if they are not already.
- Floor the standard deviation and add an acceptable pause sized to your worst normal garbage collection.
- Put a lease or fencing token behind every action that must not happen twice, such as promotion.
- Run pause, latency and one-way loss experiments, and record TD, TMR and TM for each.