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.

Advertisement

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 sizePush (mean, range)Pull (mean, range)Push-pull (mean, range)
1,00018.1 (16 to 21)13.4 (12 to 15)9.1 (8 to 10)
10,00023.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.

Advertisement

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

One gossip round between node A and node B: digests, deltas, deltasNode A (initiator)picks one random live peerNode B (receiver)compares digests1. SYN: one digest per known endpoint (address, generation, max version)2. ACK: states where B is newer + digests of what B wants3. ACK2: states B asked forsmall: no values, just version numbersA applies B's newer states, then answers B's requestsafter ACK2 both sides hold the newer of every pairEvery node runs this once per second against a random peernews reaches all n nodes in O(log n) rounds; per-node cost stays one exchange per second
Figure 1. A Cassandra-style gossip round. The SYN carries only version numbers; data moves in the ACK and ACK2, and only where the two sides differ.

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, request

New 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

ChoiceGossip (eventually consistent)Central registry (ZooKeeper, etcd)SWIM-style probing
Single point of failureNoneThe registry quorumNone
Consistency of viewsConverges in secondsLinearizableConverges; membership only
Per-node costConstant per round, digest grows with nWatches and sessionsConstant per period
Failure detectionIndirect, via heartbeat spreadSession expiryDirect and indirect probes
Good forSoft state, liveness, metadataLeadership, locks, ownershipFast 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

  1. Run the simulation above at your cluster size and note the push-pull round count; that is your best-case propagation time in seconds.
  2. List every value your system gossips and its size. Move anything large or anything that needs agreement out of gossip.
  3. 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.
  4. Alert on clock skew between nodes, since generations depend on the wall clock.
  5. Rehearse a decommission and a node replacement in staging, and confirm the departed node's record is gone everywhere after the quarantine window.
  6. Graph per-node gossip latency and up or down flapping, and investigate GC and network saturation before tuning detector thresholds.
Key takeaway: Gossip spreads cluster state through one random exchange per node per round: push-pull measured about nine rounds for 1,000 nodes and twelve for 10,000. Single-writer records versioned by generation and counter let nodes reconcile through a compact digest, and the SYN, ACK, ACK2 exchange moves only what differs. Use gossip for liveness and soft metadata, control removal and clocks, and leave agreement to consensus.