A Cassandra cluster has no master node that knows the membership list. Every node has to learn by itself which other nodes exist, which token ranges they own, which datacenter and rack they sit in, which schema version they run, and whether they are alive. Gossip is the mechanism that spreads this information: a small, periodic, peer-to-peer exchange that every node runs against random peers, so that a fact known to one node is known to all of them within seconds.
This page covers the state exchanged and how conflicts resolve, the three-message round, peers and seeds, spread rate, failure detection, operator tools, and what Cassandra 6.0 changes.
What gossip is for, and what it is not
Gossip carries cluster metadata, not data. Writes and reads travel over the same internode messaging port but in their own messages; gossip never carries a row. What it carries is small: for each node, a few dozen versioned key-value pairs describing that node.
That metadata feeds almost every other subsystem. The snitch learns other nodes' datacenter and rack from it when you use GossipingPropertyFileSnitch. The token map that routes every request is built from the tokens nodes gossip; see the token ring and vnodes. A joining node announces itself through gossip before it streams its ranges. And hint replay starts when gossip reports a node back up.
Gossip is eventually consistent: two nodes can briefly disagree about a third. For liveness that is fine; for token ownership and schema it has caused real trouble, which is why Cassandra 6.0 moves those out of gossip.
The state each node carries
Each node keeps an endpoint state for every node it knows about, including itself. An endpoint state has a heartbeat part and a map of application states:
| Piece | What it holds | Who changes it |
|---|---|---|
| Generation | Set when the node process starts (a timestamp in seconds), stored locally so it increases across restarts | The owning node, once per start |
| Heartbeat version | A counter the owner bumps about once a second | The owning node |
| STATUS / STATUS_WITH_PORT | Lifecycle state and tokens: BOOT, NORMAL, LEAVING, LEFT, shutdown and others | The owning node, or an operator command acting for a dead one |
| LOAD | Data size on disk | The owning node, periodically |
| SCHEMA | A hash of the node's schema version | The owning node, after a schema change |
| DC, RACK | Location, from the snitch | The owning node |
| RELEASE_VERSION, HOST_ID, TOKENS, addresses | Identity and how to reach the node | The owning node |
Every value carries a version from a single counter on the owning node, so each change gets a higher version than everything before it. The generation separates one run of the process from the next: after a restart it is higher, so peers replace everything they held about the old run even though the version counter starts low again.
So conflict resolution is trivial: for any endpoint, the higher (generation, version) wins, and only the owner creates new versions. The Cassandra documentation calls this a vector clock of (generation, version) tuples.
One round: SYN, ACK, ACK2
Once a second, each node bumps its own heartbeat and starts a three-message round with a chosen peer:
- SYN. The initiator sends a digest for every endpoint it knows: address, generation and the highest version it holds. Digests are small, so the message stays compact even for large clusters.
- ACK. The receiver compares the digests with its own state. For endpoints where it holds newer state, it sends that state. For endpoints where the initiator is ahead, it sends back a digest asking for the newer data.
- ACK2. The initiator sends the states that were requested.
Afterwards both nodes hold the newer of the two views for every endpoint, and each newer heartbeat is recorded as an arrival in the local failure detector, so liveness data is gathered without any direct ping.
Choosing peers, and what seeds are for
The architecture documentation describes each round as: gossip with a random live node; then, with some probability, with a node currently believed unreachable, so that a recovered node is noticed; then with a seed if the random choice did not already pick one. The unreachable probe is how a partitioned node gets reconnected without an operator.
Seeds are ordinary nodes listed in every node's seed_provider setting in cassandra.yaml. A starting node contacts them first to learn the cluster, and the extra gossip they receive keeps the cluster from splitting into islands. They are not leaders and hold no extra data. A node that lists itself as a seed does not bootstrap, so never put a new node in its own seed list.
Use two or three stable seeds per datacenter and the same list on every node; making every node a seed only adds traffic.
How fast news spreads: a worked simulation
With every node starting a round each second, a fact spreads like an epidemic: the number of nodes that know it roughly doubles each round. The simulation starts with one informed node and counts rounds until all know, 20 runs per cluster size:
import math
import random
def rounds_to_converge(n, seed):
"""Push-pull gossip: each round, every node exchanges state with one random peer."""
rng = random.Random(seed)
knows = [False] * n
knows[0] = True # node 0 has a new STATUS value
rounds = 0
while not all(knows):
rounds += 1
snapshot = knows[:]
for node in range(n):
peer = rng.randrange(n - 1)
peer += peer >= node # any node but itself
if snapshot[node] or snapshot[peer]:
knows[node] = knows[peer] = True
return rounds
for n in (10, 100, 1000, 5000):
runs = [rounds_to_converge(n, seed) for seed in range(20)]
print(f"{n:5d} nodes: mean {sum(runs) / len(runs):4.1f} rounds, "
f"worst {max(runs):2d}, log2(n) = {math.log2(n):4.1f}")Running it prints:
10 nodes: mean 3.2 rounds, worst 4, log2(n) = 3.3
100 nodes: mean 6.5 rounds, worst 7, log2(n) = 6.6
1000 nodes: mean 9.1 rounds, worst 10, log2(n) = 10.0
5000 nodes: mean 10.8 rounds, worst 12, log2(n) = 12.3The model is simplified, but the scaling holds: spread time grows with the logarithm of cluster size. A 1,000-node cluster learns a change in about ten rounds, roughly ten seconds, and five times more nodes adds one or two rounds. What limits very large clusters is message size, which grows with the number of endpoints.
Failure detection is local: phi accrual
Gossip spreads heartbeats, but no node gossips the opinion that another node is down. Each node runs its own failure detector over the heartbeat arrivals it has seen for each peer and decides locally. Two nodes can therefore disagree about whether a third is up, and nodetool status reports only the opinion of the node you ran it on.
Cassandra uses a variant of the phi accrual failure detector. Instead of a fixed timeout, it keeps a window of recent inter-arrival times per peer and computes phi: how unlikely, given past arrivals, it is to have heard nothing for this long. Above phi_convict_threshold (8 by default in cassandra.yaml) the peer is marked down. Under an exponential model of arrivals, phi is the silence divided by the mean interval times ln 10:
import math
def phi(silence_s, mean_interval_s):
"""Phi under an exponential model of heartbeat inter-arrival times.
A model, not Cassandra's exact detector, which seeds and bounds its
arrival window; use it to reason about direction and scale.
"""
# P(next heartbeat later than t) = exp(-t / mean); phi = -log10 of that.
return silence_s / (mean_interval_s * math.log(10))
def seconds_to_convict(threshold, mean_interval_s):
return threshold * mean_interval_s * math.log(10)
for mean in (1.0, 1.5, 3.0):
print(f"mean interval {mean:.1f}s: phi after 5s silence = {phi(5, mean):4.2f}; "
f"convicted at 8 after {seconds_to_convict(8, mean):4.1f}s, "
f"at 12 after {seconds_to_convict(12, mean):4.1f}s")Running it prints:
mean interval 1.0s: phi after 5s silence = 2.17; convicted at 8 after 18.4s, at 12 after 27.6s
mean interval 1.5s: phi after 5s silence = 1.45; convicted at 8 after 27.6s, at 12 after 41.4s
mean interval 3.0s: phi after 5s silence = 0.72; convicted at 8 after 55.3s, at 12 after 82.9sThese are model outputs, not measured Cassandra conviction times, but they show how the knob behaves. Raising the threshold from 8 to 12 makes detection about 50 percent slower, and the detector adapts: when heartbeats arrive less regularly, as over a congested cross-region link, conviction slows with them. Raise phi_convict_threshold to 10 or 12 only where nodes flap without real failures, since it also delays rerouting away from a dead node. A down node stops receiving reads and starts accumulating hints, but it still owns its ranges until an operator removes it.
Lifecycle states, generations and restarts
A node's STATUS (STATUS_WITH_PORT since 4.0) records where it is in its life. A joining node gossips BOOT with its tokens, which nodetool status shows as UJ, streams its data, then gossips NORMAL. A node being decommissioned gossips LEAVING, streams its ranges away, then LEFT. A node replacing a dead one announces that it is taking over the dead node's tokens. A clean shutdown gossips a shutdown state so peers mark it down at once instead of waiting for the detector.
Departure states such as LEFT linger for a while, and recently removed endpoints are quarantined, so a node that missed the message cannot reintroduce the departed node from stale state. The delays derive from the ring delay setting and vary by version; check your version's source before relying on a number.
On restart a node comes back with a higher generation and peers discard its old state wholesale. Cassandra stores the last generation locally and never starts with a lower one, even if the clock moved backwards, or peers would ignore the new run as stale.
Reading gossip as an operator
# Who is up, from THIS node's point of view
nodetool status
# Do all nodes agree on the schema? More than one version listed = disagreement.
nodetool describecluster
# Raw gossip state: generation, heartbeat and versioned application states
nodetool gossipinfo
# Last resort for a dead node that removenode cannot clear (see below)
nodetool assassinate 10.0.1.40gossipinfo shows the raw endpoint states this node holds. Each application state is printed as name, version and value:
$ nodetool gossipinfo # abbreviated and illustrative; fields vary by version
/10.0.1.12:7000
generation:1727850012
heartbeat:48211
STATUS_WITH_PORT:21:NORMAL,-1844674407370955162
LOAD:48190:2.1489632E11
SCHEMA:15:5a3c0e1e-...
DC:9:dc1
RACK:11:rack1
RELEASE_VERSION:5:5.0.x
HOST_ID:3:6f1b...
/10.0.1.13:7000
generation:1727851200
heartbeat:47002
STATUS_WITH_PORT:21:NORMAL,-922337203685477581
...A heartbeat that stops increasing between two calls means this node has not heard from that peer. Different SCHEMA values mean a schema change has not converged. A LEFT entry that never goes away is a ghost. For a dead node, use nodetool removenode with its host ID first; assassinate drops gossip state without streaming any data, so keep it as a last resort and repair afterwards.
What Cassandra 6.0 changes
Cassandra 6.0 introduces Transactional Cluster Metadata (CEP-21). Membership, token ownership and schema move from gossip into a distributed, linearized log managed by a Cluster Metadata Service, so every node applies the same metadata changes in the same order. This removes a class of problems that came from gossip's eventual consistency, such as concurrent topology changes producing inconsistent ownership views or schema disagreements that needed manual cleanup. Gossip remains part of the system, but it no longer decides who owns which tokens or which schema is current.
Moving to it is an explicit, version-specific migration step; read the 6.0 upgrade guide first. On 5.0 and earlier, everything above is how your cluster works today.
Failure modes
- Flapping nodes. Nodes marked down and up repeatedly, with hints piling up, often caused by long garbage-collection pauses or a saturated network rather than real failures. Fix the pause first; raise
phi_convict_thresholdonly after. - Schema disagreement. More than one schema version in
describeclusterafter a change, often from concurrent DDL from several clients. Run schema changes from one place, one at a time, and wait for agreement before the next. - Ghost endpoints. A removed node still shown by some nodes. Use
removenode, thenassassinateas a last resort. - Seed misconfiguration. Different seed lists across nodes or across datacenters can leave a datacenter gossiping only within itself. Keep one list everywhere, with seeds in every datacenter.
- Firewall gaps. If the internode port (7000 by default) is closed in one direction, gossip half-works and nodes disagree about liveness. Test connectivity both ways between every pair of racks and datacenters.
Trade-offs
| Property | What gossip gives | What it costs |
|---|---|---|
| No coordinator | No single point of failure for metadata | No single source of truth at any instant |
| Random peers, 1 s rounds | Spread in about log2(n) seconds | Messages grow with the number of nodes |
| Local failure detection | Adapts to each link's behaviour | Nodes can disagree about who is down |
| Eventual consistency for topology | Simple, partition-tolerant | Topology and schema races, addressed in 6.0 |
What to do next
- Check that every node has the same seed list, with two or three stable seeds per datacenter.
- Run
nodetool describeclusterand confirm there is exactly one schema version. - Run
nodetool gossipinfoon two nodes and compare heartbeats and statuses for a few peers. - Search your logs for nodes being marked down and up; correlate any flapping with GC pause logs before touching the failure detector.
- Document a dead-node procedure: replace or
removenodefirst,assassinateonly as a last resort, repair afterwards. - Verify the internode port is open in both directions between all racks and datacenters.
- If 6.0 is on your roadmap, read the Transactional Cluster Metadata upgrade notes and plan the migration step.