SysDesignPrep.com
Study guide 80 of 183

Gossip protocols and failure detection

How large clusters know who is alive without a central coordinator: heartbeats and timeouts, phi accrual failure detectors, gossip dissemination and SWIM, membership changes, split brain, and how these show up in databases, caches and service discovery.

Reading is half of it. See this used in a real interview: walk through Design a Distributed Key-Value Store →

In a cluster of hundreds or thousands of nodes, every node needs a reasonably current view of which others are alive and what they own. A central coordinator that everyone pings becomes a bottleneck and a single point of failure. Failure detectors decide when a node is probably dead; gossip protocols spread that knowledge through the cluster like a rumour. Cassandra, Consul, Redis Cluster and many others use them, and they appear in key-value store and distributed cache designs.

You cannot know for sure

Over a network, a dead node, a slow node and a broken link look the same: no response. Every failure detector is a guess with two kinds of error:

  • False positives: declaring a slow but alive node dead, which triggers needless failovers and data movement.
  • Slow detection: taking too long to notice a real failure, while requests keep failing.

Shorter timeouts detect faster but produce more false positives, especially during garbage collection pauses or network congestion.

Heartbeats and timeouts

The simplest detector: each node sends a heartbeat every second; if none arrives for, say, 5 seconds, suspect it. Easy, but a fixed timeout ignores how variable the network actually is.

Phi accrual failure detector

Instead of yes or no, the phi accrual detector (used by Cassandra and Akka) outputs a suspicion level based on the observed distribution of heartbeat intervals. If heartbeats usually arrive every second with little variance, a 3-second gap is very suspicious; on a jittery network, less so. Applications pick a phi threshold matching how aggressive they want to be. It adapts automatically to network conditions.

Gossip

In a gossip (epidemic) protocol, every node periodically (say every second) picks a few random peers and exchanges state: membership lists, heartbeat counters, versions of metadata. Information spreads exponentially: a change reaches all N nodes in about log(N) rounds, so a thousand-node cluster converges in roughly ten rounds.

Properties:

  • Scalable: each node does constant work per round regardless of cluster size.
  • Robust: no central point; messages lost or nodes failing just slow convergence slightly.
  • Eventually consistent: nodes may briefly disagree about membership.

Heartbeat counters gossiped this way also feed failure detection: if a node's counter stops increasing everywhere, it is probably dead.

SWIM

SWIM (used by Consul and Serf, via the memberlist library) improves on plain heartbeating:

  1. Each round, a node pings one random member directly.
  2. If there is no reply, it asks a few other members to ping that node indirectly, which avoids false positives caused by one bad link.
  3. Still no reply: the node is marked suspect, and it gets a chance to refute the suspicion before being declared dead.
  4. Membership updates piggyback on ping messages, so dissemination is nearly free.

Network load per node stays constant, and detection time is predictable.

Membership changes

Joining nodes contact seed nodes, then learn the rest through gossip. Ownership changes (new token ranges on a hash ring) also spread by gossip, with versioned state so newer information wins. Removing a node permanently is usually an explicit administrative action, separate from "currently unreachable", to avoid moving huge amounts of data because of a short outage. See consistent hashing.

Split brain

A network partition can make each side think the other is dead. With gossip-based membership, both sides keep serving, which is fine for an AP system that reconciles later (see quorums and leaderless replication) but dangerous for anything requiring a single leader. Systems that need one decision-maker use consensus with majority quorums, so the minority side cannot act. See consensus and coordination and CAP theorem.

Where you see them

  • Cassandra and ScyllaDB: gossip for membership and ring state, phi accrual for failure detection.
  • Redis Cluster: nodes gossip and vote to fail over a primary. See Design a Distributed Cache.
  • Consul and Serf: SWIM-based membership for service discovery. See service discovery and service mesh.
  • Kafka and etcd instead rely on consensus-based controllers for metadata, trading gossip's scalability for strict agreement. See how Kafka works.

In the interview

For Design a Key-Value Store: "Nodes discover membership and ring ownership through gossip, which converges in about log N rounds; failure detection uses a phi accrual detector so short pauses do not trigger failovers; temporarily unreachable nodes get hinted writes, and permanent removal is an explicit operation."

Checklist

  • Failure detection as a tunable trade-off between speed and false positives.
  • Adaptive detectors (phi accrual) or SWIM with indirect probes and suspicion.
  • Gossip for membership and metadata, converging in about log N rounds.
  • Versioned state so newer information wins.
  • Temporary unreachability separated from permanent removal.
  • Consensus, not gossip, where a single leader or strict agreement is required.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.