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.

Advertisement

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.

MetricMeaningWhat moves it
Detection time TDTime from a real crash until the detector permanently suspects the processHeartbeat interval plus the threshold or timeout
Mistake recurrence time TMRAverage time between two wrong suspicions of a live processNetwork jitter, pauses, threshold height
Mistake duration TMHow long a wrong suspicion lasts before it is correctedHeartbeat 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.

Advertisement

The detector family

Four designs cover almost everything you will meet.

DesignDecision ruleStrengthWeakness
Fixed timeoutSuspect if silence exceeds DTrivial, predictable TDOne D for every link; wrong in both directions at once
Adaptive estimate (Chen)Suspect if silence exceeds the estimated next arrival plus a safety marginTracks the mean delay of each linkMargin 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 thresholdAdapts to variance; one stream, many policiesAssumes a distribution; bad under bursty pauses
SWIM probing (Das, Gupta, Motivala)Probe one random member per period; ask k helpers to probe indirectly before suspectingConstant per-node load; filters out a bad local pathDetection is probabilistic; needs suspicion and refutation logic
One observer, three detector designs over the same heartbeat streamMonitored nodeheartbeat every TgapsArrival recordermonotonic clockSliding windowlast N intervalsFixed timeoutsilence > D: suspectAdaptive estimatenext arrival + marginPhi accrual-log10 P(later than now)booleanbooleanreal numberConsumers pick their own policyrouter: phi 5 | membership: phi 8 | failover: phi 12 plus a leaseSWIM-style probing replaces passive heartbeats with active checksProberping targetno ackk helpersping-req targetstill no ackSuspect, then confirmgossip; target can refuteDetectors only ever produce suspicion; safety comes from leases, fencing and quorums
The same heartbeat stream can feed a boolean timeout, an adaptive estimate or a continuous phi value; SWIM replaces passive heartbeats with direct and indirect probes. None of them removes the need for leases or quorums.

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 heartbeatphi (normal, sd 100 ms)
1,000 ms0.30
1,200 ms1.64
1,400 ms4.74
1,500 ms7.30
1,600 ms10.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 timeout

Indirect 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

  1. Write down, per consumer, what happens on a suspicion and what a false positive costs.
  2. 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.
  3. Compute the detection time your current threshold implies with the code above, for your real mean and standard deviation.
  4. Move heartbeats onto a dedicated thread and use a monotonic clock if they are not already.
  5. Floor the standard deviation and add an acceptable pause sized to your worst normal garbage collection.
  6. Put a lease or fencing token behind every action that must not happen twice, such as promotion.
  7. Run pause, latency and one-way loss experiments, and record TD, TMR and TM for each.
Key takeaway: A failure detector can only suspect, never know, and theory proves consensus survives that as long as the detector is eventually right and safety rests on quorums. Judge detectors by detection time, mistake recurrence and mistake duration rather than by a timeout. Phi accrual turns silence into a calibrated, per-link suspicion level that each consumer thresholds differently, but the normal and exponential models give very different detection times for the same threshold. At scale, SWIM's indirect probes keep load constant and filter bad paths. Floor the variance, allow for pauses, keep heartbeats off busy threads, and put leases or fencing behind anything that must not happen twice.