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.
- 1A node cannot tell a crashed peer from a slow one. A detector can only suspect, after a timeout.
- 2A shorter timeout finds crashes sooner and raises more false alarms. Choose a point on that curve.
- 3SWIM probes one member per period, asks others before it suspects, and lets a suspect refute. Load per node stays flat.
- 4A detector never decides ownership alone. Quorums, leases and fencing tokens stop a node that was wrongly declared dead.
- 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.
| method | where | status |
|---|---|---|
| Heartbeats to every node, fixed timeout | any | Small clusters |
| Gossip heartbeat counters, fixed timeout | gossip | Messages grow with N |
| Phi accrual on heartbeats | Cassandra | Approved |
| SWIM: ping, indirect ping, suspicion | memberlist · Serf | Approved |
| Sessions or leases in a consensus store | ZooKeeper · etcd | Approved |
| Load balancer health checks | load balancer | Routing only |
| A broken TCP connection as the signal | any | Not 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.
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.
suspected by a node marked dead by a node crashed, not noticed yet (node 1 at 20.0 s)
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.
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.
| tool | capability | what it gives a design | also used for |
|---|---|---|---|
| memberlist · Serf | SWIM with suspicion and incarnation numbers; gossip to 3 nodes every 200 ms | Membership and failure detection at a flat load per node. | Consul agents, event broadcast |
| Consul | A LAN gossip pool per data centre, a WAN pool between them | Fast detection inside a site; slower, separate traffic across sites. | Service catalog, health checks |
| Cassandra | Gossiped heartbeat state; a phi accrual detector on every node | Each node decides alone if a peer is up, with a threshold, not a fixed timeout. | Request routing to replicas |
| ZooKeeper | Sessions with a timeout; ephemeral znodes; watches | A member's znode vanishes when its session expires, and watchers hear about it. | Leader election, locks |
| etcd | Leases with keepalive; keys that end with their lease | A membership registry that a majority agrees on. | Leader keys, service discovery |
| Redis | Cluster bus: nodes ping each other; a majority of masters turns a suspicion into a failure | A failed master is replaced by one of its replicas, with no outside coordinator. | Slot ownership, config epochs |
| Postgres | No membership protocol; a failover manager renews a leader key with a time to live | A 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.
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- 1Each message carries one entry per node, so message size grows with the cluster.
- 2Local time only. Nodes never compare clocks.
- 3The one knob. It sets both the detection time and the false alarm rate.
Tested source Go: merge and check
// 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 fanout | found by all | false deaths per node-hour |
|---|---|---|
| 1 | 2.2 s | 4,800 |
| 3 | 1.3 s | 1.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.
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- 1A sliding window, so the detector follows the link as it changes.
- 2The chance that a heartbeat is still on its way, from a normal curve fitted to the gaps.
- 3Phi 8: a chance of 1 in 10^8 that the node is fine. Cassandra uses 8 by default.
Tested source Go: the detector
// 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))
}
| link | detector | detect | false per hour |
|---|---|---|---|
| steady | fixed 200 ms | 152 ms | 0 |
| steady | fixed 1 s | 952 ms | 0 |
| steady | phi 8 | 80 ms | 0 |
| noisy | fixed 200 ms | 179 ms | 633 |
| noisy | fixed 1 s | 979 ms | 0 |
| noisy | phi 8 | 303 ms | 2 |
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.
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- 1Round robin over a shuffled list: every member is probed within one pass of the list.
- 2Indirect probes go around a bad link between A and B.
- 3Only B raises its own incarnation. A higher number overrides any older news about B.
- 4The suspicion timeout. memberlist scales it with log N.
Tested source Go: one probe period · Go: apply a membership update
// 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})
}
})
}
// 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.
| setting | found by all | false suspicions per node-hour |
|---|---|---|
| SWIM, 0 indirect probes | 2.8 s | 590 |
| SWIM, 1 indirect probe | 2.7 s | 226 |
| SWIM, 3 indirect probes | 2.8 s | 31 |
16 nodes, 10% loss, 1 s suspicion timeout.
| nodes | detector | first finds it | entries per message |
|---|---|---|---|
| 16 | gossip heartbeats | 1.9 s | 16 |
| 128 | gossip heartbeats | 2.0 s | 128 |
| 16 | SWIM | 3.7 s | 0.3 |
| 128 | SWIM | 3.4 s | 0.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.
- 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.
| guard | it stops | it needs |
|---|---|---|
| Majority quorum | A minority side electing its own leader. | 2f + 1 voters |
| Lease | Two holders at once, if the old one stops in time. | Bounded clock drift and a safety margin |
| Fencing token | A paused old holder writing after the change. | Storage that checks the token |
| Power off the old node | Any 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.
| kind | what you see | what handles it |
|---|---|---|
| Crash | Silence, for good. | Any detector; confirm, then remove. |
| Graceful leave | The node says it leaves. | Remove at once. No timeout, no false alarm. |
| Crash and restart | The node is back with old news about itself. | A higher incarnation number overrides the old news. |
| Lossy link | Some heartbeats lost. | Indirect probes, suspicion, phi accrual. |
| One-way partition | A hears B; B does not hear A. | Indirect probes; a quorum before any role change. |
| Overloaded prober | It suspects healthy nodes. | Others refute; memberlist's Lifeguard makes a slow prober wait longer. |
| Gray failure | Pings 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.
| event | result | why it is safe | saved by |
|---|---|---|---|
| A node crashes | SWIM 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 packets | False suspicions rise. | Live nodes refute within the timeout, so false deaths stay rare. | Suspicion |
| A node pauses for longer than the timeout | It 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 2 | Each 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 requests | Membership says alive. | Request health checks eject it from the load balancer. | Outlier ejection |
| Many nodes restart together | Many 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 mistake | It hears the news and refutes it. | A higher incarnation number brings it back; members stop sending it work only in between. | Incarnation |
| step | add | it handles | move up when you see |
|---|---|---|---|
| 1 | A 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. |
| 2 | Leases 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. |
| 3 | SWIM 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. |
| 4 | One 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.
0 of 10 known
Why can a failure detector not be both fast and always right?
What does the indirect ping in SWIM protect against?
Why does SWIM suspect a node before it declares it dead?
What is the incarnation number for?
Why does SWIM load stay flat as the cluster grows, while gossip heartbeats grow?
How long does news take to reach every node by gossip?
Why use phi accrual instead of a fixed timeout?
A partition splits 5 nodes into 3 and 2. Each side marks the other dead. Who may keep the leader role?
A node answers every ping but fails half of its requests. What catches it?
Why should a lease holder stop work before its lease ends by its own clock?
- 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.