System Design
X6

Membership and failure

Every node needs a list of the nodes that are alive. The list comes from heartbeats and timeouts, so it is always a guess, and sometimes a wrong one.

Not startedSaved in this browser only.
  1. 1A node cannot tell a crashed peer from a slow one. A detector can only suspect, after a timeout.
  2. 2A shorter timeout finds crashes sooner and raises more false alarms. Choose a point on that curve.
  3. 3SWIM probes one member per period, asks others before it suspects, and lets a suspect refute. Load per node stays flat.
  4. 4A detector never decides ownership alone. Quorums, leases and fencing tokens stop a node that was wrongly declared dead.
X6
    A

    Dead or slow?

    what the observer sees
    heartbeats from node B, as sentcrashed✕gone for goodpaused (GC)resumes laterlossy link✕✕✕✕alive, beats lostnode A seessilence: dead or slow?The observer gets the same signal in all three cases. It can only suspect, after a timeout.
    • A garbage collector pause, swap or a VM move stops a node for seconds.
    • A lossy or congested link drops heartbeats from a live node.
    • False death: a live node that the cluster marks dead.

    Silence has three causes: a crash, a pause or a lossy link. From outside they look the same, so a detector can only suspect.

    B

    Options

    how a cluster finds failed nodes
    methodwherestatus
    Heartbeats to every node, fixed timeoutanySmall clusters
    Gossip heartbeat counters, fixed timeoutgossipMessages grow with N
    Phi accrual on heartbeatsCassandraApproved
    SWIM: ping, indirect ping, suspicionmemberlist · SerfApproved
    Sessions or leases in a consensus storeZooKeeper · etcdApproved
    Load balancer health checksload balancerRouting only
    A broken TCP connection as the signalanyNot approved

    A TCP connection can stay open to a dead host until a write or a keepalive fails, which can take minutes or hours.

    For a large cluster I use SWIM-style gossip. For the few nodes that own a role, I use leases in a consensus store.

    C

    See it: detection time against false alarms

    recorded from the lab simulation
    detector
    packet loss
    suspicion timeout

    Each node pings one member a second, asks 3 others to ping it if no ack comes in 500 ms, then suspects it. A suspect that does not refute within the timeout is dead.

    0.11101001k10k00 s2 s4 s6 s8 s10 sfalse deaths per node-hourGossip heartbeatsSWIMtime until the first node marks the crash dead500 ms1 s2 s4 s8 s
    n1n2n3n4n5n6n7n8n9n10n11n12n13n14n15n160 s10 s20 s30 s40 s

    suspected by a node marked dead by a node crashed, not noticed yet (node 1 at 20.0 s)

    A crash is found by the first node after 2.9 s and by every node after 3.1 s.False deaths: 0 in 6.01 node-hours, or 0 per node-hour. Suspicions of live nodes: 4 per node-hour.Load: 4.2 messages per node per second, 0.3 member entries each.

    16 nodes, about 2 ms one way, one crash per run. 40 seeded runs of 40 s for each setting: about 6 node-hours. A dot on the zero line had no false death in those hours.

    I choose the timeout from this curve: shorter finds crashes sooner, longer avoids marking live nodes dead. Suspicion lets SWIM have both.

    D

    SWIM, step by step

    click a step; its path lights up
    Node Aprobes one per periodOther membershear and pass gossip3 helpersping B for ANode Bthe target

    Step 1: Ping

    • Each period, node A pings the next member on its shuffled list: here B.
    • B answers with an ack. Most periods end here.

    If it fails

    No ack within 500 ms: the ping, the ack or B itself is lost. Go to step 2.

    SWIM pings one member per period, asks three others if the ping fails, and only then suspects. A live node refutes with a higher incarnation.

    E

    Capabilities used

    what each tool gives you
    toolcapabilitywhat it gives a designalso used for
    memberlist · SerfSWIM with suspicion and incarnation numbers; gossip to 3 nodes every 200 msMembership and failure detection at a flat load per node.Consul agents, event broadcast
    ConsulA LAN gossip pool per data centre, a WAN pool between themFast detection inside a site; slower, separate traffic across sites.Service catalog, health checks
    CassandraGossiped heartbeat state; a phi accrual detector on every nodeEach node decides alone if a peer is up, with a threshold, not a fixed timeout.Request routing to replicas
    ZooKeeperSessions with a timeout; ephemeral znodes; watchesA member's znode vanishes when its session expires, and watchers hear about it.Leader election, locks
    etcdLeases with keepalive; keys that end with their leaseA membership registry that a majority agrees on.Leader keys, service discovery
    RedisCluster bus: nodes ping each other; a majority of masters turns a suspicion into a failureA failed master is replaced by one of its replicas, with no outside coordinator.Slot ownership, config epochs
    PostgresNo membership protocol; a failover manager renews a leader key with a time to liveA standby is promoted only after the old leader's key expires.Planned switchover

    memberlist defaults: probe every 1 s, ack timeout 500 ms, 3 indirect probes, gossip to 3 nodes every 200 ms. Cassandra marks a host down at phi 8 by default.

    Consul and Serf give me SWIM membership. ZooKeeper sessions and etcd leases give me a membership list that a majority agrees on.

    F

    Heartbeats and a timeout

    pseudo code, gossip heartbeats
    gossip heartbeatspseudo code
    every 200 ms, on each node:
      my counter = my counter + 1
      send my whole table to 3 random nodes1
    
    on a table from a peer:
      FOR EACH node j: IF its counter for j > mine:
        keep it;  fresh[j] = now2
    
    every 50 ms, on each node:
      FOR EACH node j:
        IF now - fresh[j] > timeout3: mark j dead
    1. 1Each message carries one entry per node, so message size grows with the cluster.
    2. 2Local time only. Nodes never compare clocks.
    3. 3The one knob. It sets both the detection time and the false alarm rate.
    Tested source Go: merge and check
    Go: merge and checkgo
    // mergeTable keeps the higher heartbeat counter for each node. A counter that grew proves the
    // node was alive recently, so a node marked dead comes back.
    func (s *sim) mergeTable(i int, table []int64) {
      nd := s.nodes[i]
      for j, b := range table {
        if j != i && b > nd.view[j].beat {
          nd.view[j].beat, nd.view[j].fresh = b, s.now
          s.setState(i, j, alive)
        }
      }
    }
    
    // check marks dead every node whose counter has not grown for Timeout.
    func (s *sim) check(i int) {
      nd := s.nodes[i]
      if !nd.up {
        return
      }
      s.at(s.now+50, func() { s.check(i) })
      for j, m := range nd.view {
        if j != i && m.st == alive && s.now-m.fresh > s.cfg.Timeout {
          s.setState(i, j, dead)
        }
      }
    }
    
    gossip fanoutfound by allfalse deaths per node-hour
    12.2 s4,800
    31.3 s1.8

    Lab: 16 nodes, 5% loss, 1 s timeout. With fanout 1, fresh counters often take longer than the timeout to arrive.

    Each node bumps a counter and gossips its table. A counter that stops growing for the timeout means the node is dead.

    G

    Phi accrual

    a threshold, not a timeout
    phi accrualpseudo code
    on a heartbeat at time t:
      keep the gap t - last;  last = t
      mean, spread = of the last 1,000 gaps1
    
    phi(now) = -log10( P(gap > now - last)2 )
    suspect WHEN phi(now) >= 83
    1. 1A sliding window, so the detector follows the link as it changes.
    2. 2The chance that a heartbeat is still on its way, from a normal curve fitted to the gaps.
    3. 3Phi 8: a chance of 1 in 10^8 that the node is fine. Cassandra uses 8 by default.
    Tested source Go: the detector
    Go: the detectorgo
    // Phi is the phi accrual detector: it keeps the recent gaps between heartbeats and turns the
    // silence since the last one into a suspicion level. phi = 8 means the chance that a heartbeat
    // is still on its way is 10^-8.
    type Phi struct {
      gaps   []float64
      max    int
      minStd float64
      last   float64
    }
    
    // NewPhi keeps up to max gaps, and never lets the spread fall below minStd ms.
    func NewPhi(max int, minStd float64) *Phi { return &Phi{max: max, minStd: minStd, last: -1} }
    
    // Heartbeat records an arrival at time t.
    func (p *Phi) Heartbeat(t float64) {
      if p.last >= 0 {
        p.gaps = append(p.gaps, t-p.last)
        if len(p.gaps) > p.max {
          p.gaps = p.gaps[1:]
        }
      }
      p.last = t
    }
    
    func (p *Phi) stats() (mean, std float64) {
      for _, g := range p.gaps {
        mean += g
      }
      mean /= float64(len(p.gaps))
      for _, g := range p.gaps {
        std += (g - mean) * (g - mean)
      }
      return mean, max(math.Sqrt(std/float64(len(p.gaps))), p.minStd)
    }
    
    // Level is phi at time now: -log10 of the chance that the gap grows this long, for a normal
    // distribution fitted to the recent gaps.
    func (p *Phi) Level(now float64) float64 {
      mean, std := p.stats()
      later := 0.5 * math.Erfc((now-p.last-mean)/(std*math.Sqrt2))
      return -math.Log10(later)
    }
    
    // After is how long after the last heartbeat phi reaches threshold.
    func (p *Phi) After(threshold float64) float64 {
      mean, std := p.stats()
      return mean + std*math.Sqrt2*math.Erfcinv(2*math.Pow(10, -threshold))
    }
    
    linkdetectordetectfalse per hour
    steadyfixed 200 ms152 ms0
    steadyfixed 1 s952 ms0
    steadyphi 880 ms0
    noisyfixed 200 ms179 ms633
    noisyfixed 1 s979 ms0
    noisyphi 8303 ms2

    Lab: a heartbeat every 100 ms for 1 hour per link. Steady: 0.5 ms jitter, no loss. Noisy: 30 ms jitter, 1% loss. Detect: median of 200 crashes.

    Phi accrual learns the normal gap between heartbeats on each link and suspects when the silence becomes unlikely. One threshold fits steady and noisy links.

    H

    SWIM probes and suspicion

    pseudo code
    SWIMpseudo code
    every period (1 s), on node A:
      B = the next member in a shuffled list1
      ping B;  wait 500 ms for an ack
      no ack: ask 3 members to ping B for A2
      still no ack at the end of the period:
        suspect B
    
    on "suspect B, incarnation i":
      IF I am B: incarnation = i + 13
                 gossip "alive B, i + 1"
      ELSE: mark B suspect; start the timeout4
    the timeout ends, B still suspect:
      gossip "dead B"
    on "alive B, j" with j > the known one:
      B is alive again
    1. 1Round robin over a shuffled list: every member is probed within one pass of the list.
    2. 2Indirect probes go around a bad link between A and B.
    3. 3Only B raises its own incarnation. A higher number overrides any older news about B.
    4. 4The suspicion timeout. memberlist scales it with log N.
    Tested source Go: one probe period · Go: apply a membership update
    Go: one probe periodgo
    // probe runs once per Period on node i: ping one member, then ask K others, then suspect.
    func (s *sim) probe(i int) {
      nd := s.nodes[i]
      if !nd.up {
        return
      }
      s.at(s.now+s.cfg.Period, func() { s.probe(i) })
      t := s.nextTarget(i)
      if t < 0 {
        return
      }
      acked := false
      ack := func() { acked = true }
      s.send(i, t, 0, func() { s.send(t, i, 0, ack) }) // ping, ack
      s.at(s.now+s.cfg.ProbeTimeout, func() {
        if acked || !nd.up {
          return
        }
        for _, h := range s.peers(i, s.cfg.K, t, false) { // indirect: h pings t for i
          s.send(i, h, 0, func() {
            s.send(h, t, 0, func() {
              s.send(t, h, 0, func() { s.send(h, i, 0, ack) })
            })
          })
        }
      })
      s.at(s.now+s.cfg.Period-1, func() {
        if !acked && nd.up {
          s.apply(i, update{st: suspect, node: t, inc: nd.view[t].inc})
        }
      })
    }
    
    Go: apply a membership updatego
    // apply handles one membership update at node i: a local finding or a gossiped one.
    func (s *sim) apply(i int, u update) {
      nd := s.nodes[i]
      if u.node == i {
        if u.st != alive && u.inc >= nd.inc {
          nd.inc = u.inc + 1 // refute: I am alive, with a higher incarnation
          s.enqueue(i, update{st: alive, node: i, inc: nd.inc})
        }
        return
      }
      m := &nd.view[u.node]
      switch u.st {
      case alive:
        if u.inc <= m.inc {
          return
        }
        m.inc = u.inc
        m.token++
        s.setState(i, u.node, alive)
      case suspect:
        if m.st == dead || u.inc < m.inc || (m.st == suspect && u.inc == m.inc) {
          return
        }
        m.inc = u.inc
        m.token++
        tok := m.token
        s.setState(i, u.node, suspect)
        s.at(s.now+s.cfg.Timeout, func() {
          if nd.up && m.st == suspect && m.token == tok {
            s.apply(i, update{st: dead, node: u.node, inc: m.inc})
          }
        })
      case dead:
        if m.st == dead || u.inc < m.inc {
          return
        }
        m.inc = u.inc
        s.setState(i, u.node, dead)
      }
      s.enqueue(i, u)
    }
    

    A node pings one member per period, asks three others when the ping fails, and suspects only after that. A suspect refutes with a higher incarnation.

    I

    The knobs

    measured in the lab
    settingfound by allfalse suspicions per node-hour
    SWIM, 0 indirect probes2.8 s590
    SWIM, 1 indirect probe2.7 s226
    SWIM, 3 indirect probes2.8 s31

    16 nodes, 10% loss, 1 s suspicion timeout.

    nodesdetectorfirst finds itentries per message
    16gossip heartbeats1.9 s16
    128gossip heartbeats2.0 s128
    16SWIM3.7 s0.3
    128SWIM3.4 s0.4

    1% loss, 2 s timeout, 8 runs each. SWIM entries are changes only, so most messages carry none.

    Indirect probes cut false suspicions about 19 times at 10% loss. SWIM's load per node stays flat as the cluster grows.

    J

    How fast gossip spreads

    recorded from the lab simulation
    rounds until every node knows (mean of 200 runs)05101520258166425610244096nodes, log scalefanout 1: 21.5fanout 2: 12.8fanout 3: 9.664 times more nodes (64 to 4,096) adds 4 rounds at fanout 3, 10 at fanout 1.At 200 ms a round, fanout 3 reaches 4,096 nodes in 1.9 s.
    • Each round, every node that knows the news tells fanout random others.
    • The last few nodes take several extra rounds: random picks often hit nodes that already know.
    • memberlist sends each update 4 × log10(N + 1) times, rounded up, then drops it.

    Gossip reaches every node in a number of rounds that grows with log N. Fanout 3 reached 4,096 nodes in about 10 rounds.

    K

    Split brain, quorums and fencing

    when the detector is wrong
    time →node Alock servicenode Blease, token 7paused: GC, swap, VMA's lease runsexpireslease, token 8write, token 8wakes: write, token 7storage✓✕Storage keeps the highesttoken and refuses lower ones.
    guardit stopsit needs
    Majority quorumA minority side electing its own leader.2f + 1 voters
    LeaseTwo holders at once, if the old one stops in time.Bounded clock drift and a safety margin
    Fencing tokenA paused old holder writing after the change.Storage that checks the token
    Power off the old nodeAny action by the old node.Out-of-band control of the machine

    A wrong verdict is safe only if a minority cannot act, a lease ends before anyone takes over, and the storage refuses stale tokens.

    L

    Partial and gray failures

    not every failure is a crash
    kindwhat you seewhat handles it
    CrashSilence, for good.Any detector; confirm, then remove.
    Graceful leaveThe node says it leaves.Remove at once. No timeout, no false alarm.
    Crash and restartThe node is back with old news about itself.A higher incarnation number overrides the old news.
    Lossy linkSome heartbeats lost.Indirect probes, suspicion, phi accrual.
    One-way partitionA hears B; B does not hear A.Indirect probes; a quorum before any role change.
    Overloaded proberIt suspects healthy nodes.Others refute; memberlist's Lifeguard makes a slow prober wait longer.
    Gray failurePings pass; real requests fail or crawl.Request error rate and latency per node; eject outliers.

    Membership protocols catch silence. A gray failure answers pings and fails real work, so I also watch request errors and latency per node.

    M

    Failure cases

    what breaks, and what keeps it safe
    eventresultwhy it is safesaved by
    A node crashesSWIM suspects it, then marks it dead after the suspicion timeout.Work moves to other nodes once every member knows: 3.6 s in the lab at 1% loss and a 2 s timeout.SWIM
    The network drops 10% of packetsFalse suspicions rise.Live nodes refute within the timeout, so false deaths stay rare.Suspicion
    A node pauses for longer than the timeoutIt is marked dead and its role moves on.Its lease has ended, and storage refuses its old token when it wakes.Fencing token
    A partition splits 3 and 2Each side marks the other dead.Only the side with 3 of 5 can elect a leader or take a lease.Majority quorum
    A node answers pings but fails requestsMembership says alive.Request health checks eject it from the load balancer.Outlier ejection
    Many nodes restart togetherMany updates to gossip at once.Each update is sent a fixed number of times, so the gossip load has a ceiling.Retransmit limit
    A node is marked dead by mistakeIt hears the news and refutes it.A higher incarnation number brings it back; members stop sending it work only in between.Incarnation
    N

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1Static list2+ leases in etcd3+ SWIM gossip4+ pool per sitemore load →
    Member entries sent per node per second0.11101001k10kGossip, 16 nodes: 232 entries per node per secondGossip, 16 nodes232Gossip, 128 nodes: 1,907 entries per node per secondGossip, 128 nodes1,907SWIM, 16 nodes: 1 entries per node per secondSWIM, 16 nodes1SWIM, 128 nodes: 1.5 entries per node per secondSWIM, 128 nodes1.5entries per node per second, log scale
    stepaddit handlesmove up when you see
    1A static list in config, with load balancer health checks.A few nodes that change rarely. The balancer stops sending to a node that fails its checks.Nodes come and go on their own, or one node must own a role.
    2Leases or sessions in etcd or ZooKeeper.A membership list a majority agrees on, and one owner per role.Hundreds of members renew leases, and the consensus store spends its writes on renewals.
    3SWIM gossip, for example memberlist or Serf. Keep leases for roles.Flat load per node. News reaches 4,096 nodes in about 9.6 rounds.Members in several sites, and gossip across slow links causes false suspicions.
    4One pool per site, as in Consul's LAN and WAN pools.Fast detection inside a site, and separate timeouts across sites.Top of the ladder.

    Entries per node per second, measured in the lab: messages per node per second times member entries per message, at 1% loss. Gossip heartbeats grow with N; SWIM does not.

    I start with a static list and health checks. I add leases in etcd for the nodes that own roles, and gossip membership only when the cluster reaches hundreds of nodes.

    O

    Drill

    predict, then reveal

    0 of 10 known

    1. Why can a failure detector not be both fast and always right?

    2. What does the indirect ping in SWIM protect against?

    3. Why does SWIM suspect a node before it declares it dead?

    4. What is the incarnation number for?

    5. Why does SWIM load stay flat as the cluster grows, while gossip heartbeats grow?

    6. How long does news take to reach every node by gossip?

    7. Why use phi accrual instead of a fixed timeout?

    8. A partition splits 5 nodes into 3 and 2. Each side marks the other dead. Who may keep the leader role?

    9. A node answers every ping but fails half of its requests. What catches it?

    10. Why should a lease holder stop work before its lease ends by its own clock?

    P

    Numbers to say

    measured, derived, cited
    memberlist
    Probe every 1 s, ack in 500 ms, 3 indirect probes, gossip to 3 nodes every 200 ms.
    suspicion
    memberlist: 4 × log10(N) probe periods, at least 4. At 16 nodes, 4.8 s.
    phi
    Cassandra marks a host down at phi 8.
    detect
    Lab, SWIM, 16 nodes, 1% loss, 2 s suspicion: every node knew after 3.6 s.
    helpers
    At 10% loss: 590 false suspicions per node-hour with no indirect probes, 31 with 3.
    spread
    Gossip rounds grow with log N: 5.2 rounds for 64 nodes, 9.6 for 4,096, at fanout 3.

    Lab numbers come from a seeded simulation of 16 nodes with about 2 ms one-way delay, not from a machine. Defaults from the memberlist source and the Cassandra configuration.