Every distributed system needs a member list: which nodes exist, which are alive, and where to reach them. A load balancer needs it to route, a database needs it to place replicas, and a service mesh needs it to fail over. The obvious design, where every node heartbeats every other node, works for ten nodes and falls over at a thousand. A central registry works until the registry is the thing that fails.

Gossip-based membership solves this with constant work per node and no coordinator. This article explains the protocol most production systems use, SWIM (Das, Gupta and Motivala, 2002), and the refinements HashiCorp published as Lifeguard (Dadgar and colleagues, 2017), which are implemented in the memberlist library behind Consul, Nomad and Serf. For gossip as a general dissemination technique, see the gossip protocol overview. Here we stay on membership: how failures are detected, how the news spreads, and how to tune it.

Advertisement

The problem and its requirements

A membership protocol has two jobs: failure detection, noticing that a member has stopped responding, and dissemination, telling everyone about joins, leaves and failures. We judge it on four properties. Completeness: every failure is eventually noticed by every live member. Speed: how long that takes. Accuracy: how rarely a live member is declared dead. Scalability: how load per node and network traffic grow with cluster size n.

All-to-all heartbeating gets completeness and speed but fails on scalability. Each node sends n−1 heartbeats per interval, so the network carries on the order of n² messages, and one slow node's view of time drives false positives everywhere. SWIM's insight is to separate the two jobs. Detection uses a small, fixed amount of probing per node per period. Dissemination rides on the probe messages themselves, so it costs almost nothing extra.

SWIM failure detection: the probe cycle

Time is divided into protocol periods, for example one second. In each period every node does the following. It picks one target. Implementations walk a shuffled list in round-robin order rather than picking purely at random, which bounds how long any member can go unprobed. It sends the target a ping and waits for an ack up to a probe timeout shorter than the period. If no ack comes, it asks k other members to ping the target on its behalf (ping-req) and relay any ack. If neither the direct ping nor any of the indirect pings succeeds by the end of the period, the target becomes suspect.

The indirect step is what makes SWIM accurate. A lost packet or a congested link between A and B would make a direct-only detector accuse B wrongly. Requiring k independent paths to fail as well makes a false accusation much less likely. The load stays constant: each node sends one ping per period and occasionally k ping-reqs, whatever the value of n. Because each member is probed by someone in nearly every period, the expected time to first detection is about one period plus timeouts.

One SWIM protocol period at node A (probe target B, helpers C, D, E)Node AproberNode Brandom targetC, D, Ek = 3 helpers1. ping (with piggybacked updates)2. ack within probe timeout?3. no ack: ping-req(B)4. ping B5. relayed ackOutcome at end of periodAny ack arrived: B stays alive. Updates carried on every message spread as a side effect.No ack at all: A gossips Suspect(B, incarnation i). B has a suspicion timeout to refute.B hears it and gossips Alive(B, i+1): the higher incarnation overrides the suspicion.Nobody refutes before the timeout: Dead(B) is gossiped; B returns only with a higher incarnation.
A single protocol period. Steps 3 to 5 run only when the direct ping fails; the helpers give the target k more independent network paths before anyone suspects it.
Advertisement

Suspicion and incarnation numbers

Declaring a member dead on the first failed probe would make every GC pause an outage. SWIM adds a suspect state. A suspect member is still treated as part of the cluster, and the suspicion is gossiped. If the member hears that it is suspected, it refutes the suspicion by gossiping that it is alive with a higher incarnation number, a counter that only the member itself may increment. If no refutation arrives before the suspicion timeout, the suspecting node declares it dead.

Incarnations give every node the same rule for deciding which message wins, with no clocks involved:

Incoming messageOverrides local state when
Alive(m, i)the local incarnation for m is lower than i; this is also how a member marked dead comes back
Suspect(m, i)local is Alive with incarnation i or lower, or Suspect with incarnation lower than i
Dead(m, i)i is at least the local incarnation and m is not already dead

The asymmetry is deliberate. At equal incarnation, suspicion beats alive, so a stale alive rumour cannot erase a fresh accusation. Only the accused can bump the number. A merge function that applies these rules is short. Here it is in Python, with transport and timers left as stubs:

import random, time
from dataclasses import dataclass

ALIVE, SUSPECT, DEAD = 0, 1, 2

@dataclass
class Member:
    name: str
    addr: str
    incarnation: int
    state: int
    since: float

def apply(members: dict, me: str, msg) -> bool:
    """Merge one gossiped (name, state, incarnation). Returns True if it is news to re-gossip."""
    m = members.get(msg.name)
    if msg.name == me:
        if msg.state != ALIVE and msg.incarnation >= state_of_me.incarnation:
            state_of_me.incarnation = msg.incarnation + 1      # refute: outbid the rumour
            enqueue(Alive(me, state_of_me.incarnation))
        return False
    if m is None:
        if msg.state == DEAD:
            return False
        members[msg.name] = Member(msg.name, msg.addr, msg.incarnation, msg.state, time.time())
        return True
    newer = msg.incarnation > m.incarnation
    same = msg.incarnation == m.incarnation
    if (msg.state == ALIVE and newer) or \
       (msg.state == SUSPECT and m.state != DEAD and (newer or (same and m.state == ALIVE))) or \
       (msg.state == DEAD and m.state != DEAD and msg.incarnation >= m.incarnation):
        m.state, m.incarnation, m.since = msg.state, msg.incarnation, time.time()
        return True
    return False

def protocol_period(members, me, k=3, probe_timeout=0.5):
    target = next_in_shuffled_round_robin(members)     # every member probed once per round
    if ping(target, piggyback=take_updates()):
        return
    helpers = random.sample([m for m in alive(members) if m.name != target.name], k)
    if any_ack(ping_req(h, target) for h in helpers):
        return
    if apply(members, me, Suspect(target.name, target.incarnation)):
        enqueue(Suspect(target.name, target.incarnation))
        start_suspicion_timer(target)       # on expiry: apply and enqueue Dead(...)

Dissemination by piggybacking

SWIM sends no separate broadcast. Each node keeps a buffer of recent membership updates and attaches as many as fit to every ping, ping-req and ack it sends. A node that receives a new update applies it and puts it in its own buffer. This is epidemic spread: the number of informed nodes roughly doubles each round, so news reaches the whole cluster in a number of rounds that grows with log n.

Each update is retransmitted only a limited number of times, and the buffer prefers the updates sent least often so that fresh news is not crowded out. In memberlist the limit is RetransmitMult × ceil(log10(n + 1)). With the LAN default RetransmitMult of 4 and 1,000 nodes, ceil(log10(1001)) is 4, so each node sends each update 4 × 4 = 16 times. (The suspicion timeout below uses log10(n) without rounding up, so the two scales differ slightly.) Memberlist also runs a separate gossip timer (every 200 ms to 3 random nodes by default) so updates do not wait for the next probe. Because messages are UDP datagrams, the piggyback space is bounded by a safe packet size. A burst of joins therefore drains over several rounds instead of in one big message.

Worked example: a 1,000-node cluster on memberlist defaults

memberlist's DefaultLANConfig, read from the source at the time of writing, uses: ProbeInterval 1 s, ProbeTimeout 500 ms, IndirectChecks 3, SuspicionMult 4, SuspicionMaxTimeoutMult 6, GossipInterval 200 ms, GossipNodes 3, RetransmitMult 4, PushPullInterval 30 s and AwarenessMaxMultiplier 8. Check config.go for the version you vendor, because defaults change.

Suppose node B's process freezes. Within about one probe interval some node finds no direct ack after 500 ms and sends three ping-reqs, and those fail too. B becomes suspect, and the suspicion spreads through gossip, reaching most of the cluster within a few 200 ms gossip rounds. The minimum suspicion timeout is SuspicionMult × max(1, log10(n)) × ProbeInterval = 4 × 3 × 1 s = 12 s. The maximum is SuspicionMaxTimeoutMult times that, 72 s. Lifeguard starts at the maximum and lowers it as independent confirmations arrive, which is the subject of the next section. So a real failure is typically declared somewhere between about 12 s and a little over a minute. If B was only paused and resumes inside that window, it refutes with a higher incarnation and nothing else happens.

The cost per node is about one ping and ack per second, a few gossip packets every 200 ms, and a full-state TCP exchange with one random peer every 30 s. None of that grows with n except the size of the push-pull state, so 1,000 nodes costs each node about what 10 nodes does.

Lifeguard: when the detector is the sick one

Production experience at HashiCorp showed a failure mode the original paper does not address. A node that is itself slow, starved of CPU or stuck in a long GC pause, processes acks late, accuses healthy peers, and spreads false suspicions. Lifeguard adds three mechanisms.

  • Local health awareness. Each node keeps a health score that rises when its own probes fail oddly, for example missed nacks from helpers or having to refute suspicions of itself. Probe interval and timeout are scaled up by this score, up to AwarenessMaxMultiplier. A sick node slows down instead of accusing others.
  • Dynamic suspicion timeout. The timeout starts long and shrinks as independent members report the same suspicion, so a real failure confirmed by many peers is declared quickly, while a lone accusation waits.
  • Buddy system. When a node probes a member it already suspects, it tells that member directly, so the accused can refute without waiting for the rumour to reach it.

The Lifeguard authors report a large reduction in false positives with little change in detection time for real failures. That is why memberlist turns these on by default.

Anti-entropy, and membership is not consensus

Piggybacking forgets each update after a limited number of retransmits, so a node that was partitioned away can miss news for good. Memberlist fixes this with periodic push-pull: two random nodes exchange their full member lists over TCP and merge them with the same precedence rules. This also lets a joining node learn the whole cluster from one seed. Cassandra takes a different route: nodes gossip versioned state for each endpoint, and each node decides liveness locally with a phi accrual failure detector instead of a SWIM probe cycle.

Above all, the member list is eventually consistent. Two nodes can briefly disagree about who is alive, and during a network partition each side will declare the other dead. That is fine for routing and for placing work. It is not fine for anything that needs a single answer, such as electing one leader or deciding who owns a lock. For those, put membership changes through consensus, as Consul does with Raft for its catalog (see Raft consensus), or require a quorum before acting.

Failure modes and operational guidance

  • Asymmetric reachability. If A can reach B but B cannot reach A, acks are lost in one direction. Indirect probes help, but persistent asymmetry causes flapping. Alert on the rate of refutations per node.
  • Wrong advertised address. In containers or behind NAT, a node advertises an address its peers cannot reach. It joins, then fails every probe. Set the advertise address explicitly.
  • UDP blocked or fragmented. Firewalls that allow TCP but not UDP give you push-pull without probing. Oversized packets are dropped silently. Keep messages under the path MTU.
  • Resurrected names. A node that restarts and announces an old incarnation is ignored while its peers still hold it as dead, until it announces a higher one. Memberlist also has a reclaim window for a dead name rejoining from a new address. Know how yours handles restarts.
  • Overloaded hosts. CPU starvation is the most common cause of false suspicion. Watch the local health score, and give the gossip thread priority.
  • WAN use with LAN settings. Cross-region latency breaks a 500 ms probe timeout. Use the WAN profile or a separate gossip pool per region.
  • Mixed encryption keys during rotation. Install the new key on every node, switch the primary, and only then remove the old one.

Track these metrics per node: probe failures, suspicions raised, refutations sent, members declared dead, push-pull duration and the local health score. A steady trickle of suspicions that are then refuted is the earliest warning of network or host trouble.

What to do next

  1. Decide what consumes membership in your system, and separate routing uses (gossip is fine) from ownership and leadership uses (need consensus or a quorum).
  2. Start from a proven implementation such as memberlist rather than writing SWIM from scratch, and pin the version.
  3. Write down detection time and false-positive targets, then derive the probe interval and suspicion multiplier from them.
  4. Set advertise addresses explicitly, open UDP and TCP on the gossip port, and keep packets under the MTU.
  5. Export the per-node metrics above and alert on refutation rate and health score.
  6. Test with a fault injector: kill a process, pause one with SIGSTOP, drop packets one way, and partition the network. Confirm that detection times and refutations behave as calculated.
Key takeaway: Gossip-based membership gives every node a near-current view of the cluster at constant cost per node. SWIM probes one random member per period, uses k indirect pings before suspecting it, lets the accused refute by raising its incarnation, and spreads all of this on the probe messages themselves. Lifeguard stops slow nodes from accusing healthy ones, and push-pull heals what piggybacking forgets. Treat the result as eventually consistent: use it for routing, put leadership and ownership behind consensus, tune from measured detection and false-positive targets, and test with real faults.