Every node in a cluster needs to know who else is in it, where they live, and facts about each: how loaded it is, whether it is joining, leaving or alive. A central registry can hold that list, but becomes a dependency every node must reach. Gossip takes the opposite approach. Each node periodically talks to one random peer and the two exchange what they know. News spreads like an infection, and within a few seconds every node has it.
This article explains gossip as a dissemination and reconciliation mechanism: why random exchanges converge so fast, how versioned state lets two nodes reconcile with a few bytes of digest, and how a production system, Apache Cassandra, structures its rounds. The failure-detection side, SWIM-style probing with suspicion and incarnation numbers, is covered separately in gossip-based membership with SWIM; here we focus on how state moves.
The problem gossip solves
Consider 500 database nodes, each holding a small record about itself: address, rack, a load figure and a status such as NORMAL or LEAVING. Every node needs a fresh copy of all 500 records.
Broadcast is the naive answer, but it needs every node to hold a reliable list of everyone else, the very thing we are building, and silently loses updates to nodes that were briefly unreachable. A central registry such as ZooKeeper or etcd fixes consistency but adds a strongly consistent quorum to the critical path of membership and makes a registry outage a cluster outage.
Gossip trades strong consistency for robustness. Every node is equal, per-node cost is constant per round, and a node that was cut off catches up the next time it talks to anyone. The price is that for a few seconds different nodes hold different views; anything that cannot tolerate that, such as lock ownership, needs consensus.
Epidemics from first principles: push, pull and push-pull
Model the cluster as n nodes and one new rumour that starts at a single node. Time moves in rounds; in each round every node picks one peer uniformly at random. There are three ways to use that contact. In push, an informed node sends the rumour to its peer. In pull, every node asks its peer for news and learns the rumour if the peer has it. In push-pull, the two exchange in both directions.
Push is fast at the start, when the informed set roughly doubles each round, and slow at the end, when most pushes land on nodes that already know. Pull has the opposite shape: slow while few nodes have the rumour, then every straggler finds it almost immediately. Push-pull gets the fast start of push and the fast finish of pull. Classical analyses give about log2(n) + ln(n) rounds for push, and Karp and co-authors (2000) showed push-pull finishes in about log3(n) plus a smaller order term.
The simulation below measures it rather than asserting it.
import random
def rounds_to_spread(n, mode, seed):
"""One rumour starts at node 0. Each round every node contacts one random peer."""
rng = random.Random(seed)
informed = [False] * n
informed[0] = True
count, rounds = 1, 0
while count < n:
rounds += 1
nxt = informed[:] # synchronous rounds: decide from last round's state
for i in range(n):
j = rng.randrange(n - 1)
if j >= i:
j += 1 # never pick yourself
if mode in ("push", "pushpull") and informed[i]:
nxt[j] = True
if mode in ("pull", "pushpull") and informed[j]:
nxt[i] = True
informed = nxt
count = sum(informed)
return rounds
for n in (1_000, 10_000):
for mode in ("push", "pull", "pushpull"):
rs = [rounds_to_spread(n, mode, s) for s in range(20)]
print(n, mode, sum(rs) / len(rs), min(rs), max(rs))Twenty runs per configuration produced these round counts:
| Cluster size | Push (mean, range) | Pull (mean, range) | Push-pull (mean, range) |
|---|---|---|---|
| 1,000 | 18.1 (16 to 21) | 13.4 (12 to 15) | 9.1 (8 to 10) |
| 10,000 | 23.5 (21 to 27) | 17.9 (16 to 22) | 11.7 (11 to 13) |
A tenfold larger cluster costs only a few more rounds, and push-pull roughly halves the time of push. At one round per second, a 1,000-node push-pull cluster spreads a change in about nine seconds in this idealised model.
Rumour mongering versus anti-entropy
Demers and colleagues at Xerox PARC (1987) separated two styles. Rumour mongering spreads only hot changes and stops, which is cheap but can miss a few nodes. Anti-entropy compares entire state on every contact, which eventually repairs anything but ships more data.
Membership gossip is anti-entropy made cheap by versioning: nodes compare a compact digest of version numbers and ship only the records that differ. The same idea applied to replicated data, rather than membership metadata, is covered in anti-entropy for replicas.
Versioned state: generations and versions
Two nodes must decide, without asking anyone, which copy of a record is newer. Cassandra uses two numbers per endpoint. The generation is set when the node process starts, from the wall clock in seconds, so a restarted node always has a higher generation than its previous life. The version is a counter that increases every time that node changes anything about itself, including a heartbeat that ticks every round. Each piece of application state, such as load or status, carries the version at which it was last set.
Only the owning node ever writes its own record. Other nodes just relay copies. That single-writer rule is what makes the ordering trivial: compare generations first, and if they are equal, compare versions. No vector clocks, no conflict resolution.
from dataclasses import dataclass, field
@dataclass
class EndpointState:
generation: int # set once at process start (seconds since epoch)
heartbeat: int = 0 # bumps every gossip round
app: dict = field(default_factory=dict) # key -> (value, version)
def max_version(self):
return max([self.heartbeat] + [v for _, v in self.app.values()])
class LocalNode:
def __init__(self, addr, generation):
self.addr = addr
self.version = 0
self.states = {addr: EndpointState(generation)}
def _next(self):
self.version += 1
return self.version
def tick(self): # once per round
self.states[self.addr].heartbeat = self._next()
def set(self, key, value): # e.g. ("LOAD", "412GB") or ("STATUS", "LEAVING")
self.states[self.addr].app[key] = (value, self._next())Heartbeat and application state share one counter, so the maximum version summarises everything a peer knows about that node: exactly what a digest needs.
The three-message exchange, worked through
Cassandra names the three messages GossipDigestSyn, GossipDigestAck and GossipDigestAck2. The initiator sends a digest list: for every endpoint it knows, the address, the generation and the maximum version it holds. The receiver compares each digest with its own copy. Where it holds something newer, it puts the newer state into the ACK. Where the initiator holds something newer, it puts a request digest into the ACK, telling the initiator what version it already has. The initiator applies the newer states and answers the requests in ACK2.
Work through a three-node cluster. Node A knows: A at generation 1000 version 52, B at 1000 version 40, C at 1000 version 31. Node B knows: A at 1000 version 49, B at 1000 version 44, C at 1003 version 2, because C restarted three seconds ago and told B. A sends SYN with its three digests. B compares. For A, B is behind (49 against 52), so it asks for A above version 49. For B, B is ahead (44 against 40), so it includes B's state above version 40. For C, B holds a higher generation, so it includes C's entire new state, because a new generation replaces everything. A applies the B and C states, then sends ACK2 containing A's state above version 49. After one round trip of three messages, both nodes agree on all three records.
def examine(local, digests):
"""Receiver side of a SYN. Returns (states_to_send, requests)."""
send, request = {}, []
for addr, gen, ver in digests:
mine = local.states.get(addr)
if mine is None:
request.append((addr, gen, 0)) # know nothing: ask for all
elif gen > mine.generation:
request.append((addr, gen, 0)) # peer restarted: ask for all
elif gen < mine.generation:
send[addr] = mine # we have the newer life
elif ver > mine.max_version():
request.append((addr, gen, mine.max_version()))
elif ver < mine.max_version():
send[addr] = delta(mine, above=ver) # only fields newer than ver
return send, requestNew members propagate because each initiator's digest lists every endpoint it knows, and a peer that has never heard of one requests it in full. Cassandra also ignores a remote generation more than a year ahead of local time, guarding against a clock jump.
Choosing whom to gossip with
Pure randomness leaves two gaps: unreachable nodes are never contacted again, and new nodes or partitioned islands need common meeting points.
Cassandra's gossiper runs once per second and handles both. In each round a node gossips to one randomly chosen live endpoint. It may also gossip to one unreachable endpoint, with probability equal to the number of unreachable endpoints divided by the number of live endpoints plus one, so recovery probes scale with how much of the cluster is down. If the live peer it chose was not a seed, or it knows fewer live nodes than there are seeds, it may also gossip to a seed: always when it knows no live nodes at all, otherwise with probability equal to the number of seeds divided by the number of other known endpoints. Seeds are ordinary nodes listed in configuration; their only special role is to be the known rendezvous that bootstraps a new node and stitches islands back together.
Note that in Cassandra's trunk source, states such as tokens, rack and schema are now marked as derived from transactional cluster metadata (CEP-21): ownership, which needs agreement, moves to a consensus-backed log, while gossip keeps liveness and soft state. Check the release you run.
Liveness on top of gossip
Gossip itself only moves data. Failure detection is layered on top by watching the heartbeat inside each record. If node C's heartbeat version stops increasing in everyone's copy, C has probably stopped. Cassandra feeds the arrival times of heartbeat updates into a phi-accrual failure detector, which converts the gap since the last update into a suspicion level and marks the node down once it crosses a configured threshold.
A node is thus judged by whether news of its heartbeat is spreading, not by whether you can reach it. That averages out single slow links, but detection takes the spread time plus the detector's window, often several seconds.
Leaving, restarting and the zombie problem
Removing state is the hard part. If every node simply deletes decommissioned node D's record, the next SYN from a node that has not yet deleted it reintroduces D. So removal is itself state: D's record gets a terminal status such as LEFT, which gossips like any change and is kept until every copy has seen it.
A restarted node starts a new generation, so its fresh record beats every stale copy. Cassandra also quarantines a removed endpoint for a period derived from its ring delay so late gossip about it is ignored, and offers an operator command, nodetool assassinate, as a last resort for forcibly removing an endpoint whose state refuses to die. Needing it is a signal that removal did not propagate cleanly.
Failure modes and operational guidance
- Bad seed lists. If a new node's seeds are unreachable or point at another cluster, it cannot join or, worse, joins the wrong one. Use two or three seeds per data centre, keep the list identical everywhere, and do not make every node a seed, which defeats the rendezvous role.
- Clock jumps. Generations come from the wall clock. A node that starts with a clock far in the future produces a generation that peers reject or that dominates for years. Run NTP or chrony and alert on skew before it matters.
- Oversized state. Each SYN carries a digest per endpoint, so its size grows linearly with cluster size, and every round ships deltas. Do not put large values into gossip; keep it to addresses, flags and small numbers.
- Partitions. Two islands each converge on a view in which the other is down. That is correct behaviour, but anything acting on membership, such as moving data, must not treat a gossip down state as proof of death.
- Flapping under load. A node with long GC pauses gossips late and is judged late. Treat up and down flapping as a resource problem first.
Trade-offs
| Choice | Gossip (eventually consistent) | Central registry (ZooKeeper, etcd) | SWIM-style probing |
|---|---|---|---|
| Single point of failure | None | The registry quorum | None |
| Consistency of views | Converges in seconds | Linearizable | Converges; membership only |
| Per-node cost | Constant per round, digest grows with n | Watches and sessions | Constant per period |
| Failure detection | Indirect, via heartbeat spread | Session expiry | Direct and indirect probes |
| Good for | Soft state, liveness, metadata | Leadership, locks, ownership | Fast crash detection |
The pattern most mature systems reach is a blend: gossip for liveness and soft state, a consensus system for anything that needs agreement. The Dynamo paper's influence shows how far that split has spread across databases.
What to do next
- Run the simulation above at your cluster size and note the push-pull round count; that is your best-case propagation time in seconds.
- List every value your system gossips and its size. Move anything large or anything that needs agreement out of gossip.
- Confirm every node uses the same short seed list, with at least two seeds per data centre, and that seeds are not your only bootstrap path in automation.
- Alert on clock skew between nodes, since generations depend on the wall clock.
- Rehearse a decommission and a node replacement in staging, and confirm the departed node's record is gone everywhere after the quarantine window.
- Graph per-node gossip latency and up or down flapping, and investigate GC and network saturation before tuning detector thresholds.