Internode communication through gossip

Hyper-Converged Database (HCD) uses the gossip protocol to discover location and state information about other nodes in the cluster.

Gossip is a peer-to-peer communication protocol in which nodes periodically exchange state information about themselves and about other nodes they know about. The gossip process runs every second and exchanges state messages with up to three other nodes in the cluster. The nodes exchange information about themselves and about the other nodes that they have gossiped about, so all nodes quickly learn about all other nodes in the cluster. Each gossip message includes a version that allows nodes to overwrite older information with the most current state during exchanges.

Use consistent seed nodes to prevent gossip communication issues

To prevent problems in gossip communications, be sure to use the same list of seed nodes for all nodes in a cluster, with the exception that a seed node shouldn’t list itself in its own seeds list.

The seeds list is most critical the first time a node starts up, which is why you must start seed nodes before other nodes in the cluster.

By default, a node remembers other nodes it has gossiped with between subsequent restarts. This means that the seed node designation has no purpose other than bootstrapping the gossip process for new nodes joining the cluster. Seed nodes are not a single point of failure, but they don’t have any other special purpose in cluster operations beyond bootstrapping.

Don’t make every node a seed node because it increases maintenance and reduces gossip performance. Gossip optimization isn’t critical, but DataStax recommends a small seed list with two or three seed nodes per datacenter.

About failure detection and recovery

Failure detection is a method for locally determining from gossip state and history when a node in the system is down or has come back up. The database uses this information to avoid routing client requests to unreachable nodes whenever possible. The database can also avoid routing to poorly performing nodes, through dynamic snitching.

Each node’s gossip process tracks node states from two sources: directly when nodes share gossip about themselves, and indirectly when nodes share gossip received from other nodes. Rather than using a fixed threshold for marking failing nodes, the database uses an accrual detection mechanism to calculate a per-node threshold. The threshold takes into account network performance, workload, and historical conditions. During gossip exchanges, every node maintains a sliding window of inter-arrival times of gossip messages from other nodes in the cluster.

To adjust the sensitivity of the failure detector, configure the phi_convict_threshold property in the cassandra.yaml file:

Lower values increase the likelihood that an unresponsive node will be marked as down. Use the default value for most situations. In unstable for frequently congested network environments, increasing this value to 10 or 12 can reduce false failure detections. Avoid using values higher than 12 or lower than 5.

Node failures can result from various causes such as hardware failures and network outages. Node outages are often transient but can last for extended periods. Because a node outage rarely signifies a permanent departure from the cluster, it does not automatically result in permanent removal of the node from the ring. Other nodes will periodically try to re-establish contact with failed nodes to see if they are back up.

When a node comes back online after an outage, it might have missed writes for the replica data it maintains. Some missed writes can be automatically recovered through hinted handoffs. If immediate recovery is needed, or the node was down for longer than the hinted handoff window, run manual repair with nodetool repair.

If a node is dead, permanently leaving the cluster, or requires replacement, you must explicitly remove the node from the cluster to stop other nodes from trying to communicate with it and replay missed writes.

If you intend to reuse a removed node later, make sure you properly remove it from the cluster, clear the data directories, and delete stale configuration files before attempting to reuse the node. Any remnant of the prior configuration can cause the node to start improperly, resurrect stale data, or modify the cluster configuration unexpectedly.

Was this helpful?

Give Feedback

How can we improve the documentation?

© Copyright IBM Corporation 2026 | Privacy policy | Terms of use |  Manage Privacy Choices

Apache, Apache Cassandra, Cassandra, Apache Tomcat, Tomcat, Apache Lucene, Apache Solr, Apache Hadoop, Hadoop, Apache Pulsar, Pulsar, Apache Spark, Spark, Apache TinkerPop, TinkerPop, Apache Kafka and Kafka are either registered trademarks or trademarks of the Apache Software Foundation or its subsidiaries in Canada, the United States and/or other countries. Kubernetes is the registered trademark of the Linux Foundation.

General Inquiries: Contact IBM