Repair mechanisms
Over time, data in a replica can become inconsistent with other replicas due to the distributed nature of the database. Node repair processes correct inconsistencies and ensure that all nodes have the latest data. Node repair is an important part of regular maintenance for every Hyper-Converged Database (HCD) cluster.
Several repair mechanisms exist, and each serves a specific purpose in maintaining data consistency across the cluster. Database configuration settings and HCD tools are used to manage and run different types of repair.
|
DataStax recommends stopping repair operations during topology changes because repairs that run during a topology change can cause errors when repairs involve moving token ranges. |
Anti-entropy repair
Deletes and missed writes while nodes are down often cause data inconsistency.
Anti-entropy repair is an essential part of routine cluster maintenance in addition to on-demand repairs with nodetool repair.
Because data on disk is stored in immutable SSTables, the repair process requires reading SSTables to identify inconsistencies, and then rewriting the SSTables for any inconsistent data ranges. This is inherently resource intensive, especially when there are many inconsistencies involving multiple partitions. Running repairs regularly reduces the amount of data that needs to be rewritten during each repair, minimizing the impact on cluster performance.
|
Anti-entropy repair uses a structure known as a Merkle tree, which is comprised of nodes. The nodes in a Merkle tree aren’t equivalent to nodes in your HCD cluster. |
Anti-entropy repair process
The anti-entropy process does the following:
-
Build a Merkle tree for each replica.
-
Compare the Merkle trees to discover differences across replicas.
-
Update replicas to the latest version of the data.
Merkle trees are binary hash trees where the leaves are hashes of the individual key values. In HCD databases, a leaf of a Merkle tree is the hash of a row value. Each parent node higher in the tree contains a hash of its respective children. Because higher nodes in the Merkle tree represent data further down the tree, the database can check each branch independently without requiring the coordinator node to download the entire dataset.
For anti-entropy repair, HCD uses a compact tree version with a depth of 15 (215), resulting in 32,768 leaf nodes. For example, for a node containing one million partitions with one damaged partition, about 30 partitions are streamed, which is the number that fall into each of the leaves of the tree. HCD works with smaller Merkle trees because they require less storage memory and can be transferred more quickly to other nodes during the comparison process.
After the coordinator node receives the Merkle trees from the participating peer nodes, the coordinator node compares every tree to every other tree, beginning with the top node of the Merkle tree. A difference between Merkle tree nodes indicates inconsistent data on replicas for the hash range represented by the conflicting Merkle tree nodes.
If there is no difference, the data requires no repair, and the comparison process continues to the next set of Merkle tree nodes.
If there is a difference, the comparison process proceeds systematically to the left and right child nodes to determine the specific ranges of data that are inconsistent. When the conflicting range is identified, the coordinator node initiates repair on the relevant replicas, also known as the replica set. Anti-entropy repair replaces all data that corresponds to the leaves below the conflicting Merkle tree node.
Anti-entropy repair options
Building Merkle trees is resource intensive, stressing disk I/O and using memory.
The target of the nodetool repair command, the frequency at which you run repairs, and the volatility of your data determine the performance impact of an anti-entropy repair operation.
The nodetool repair command can target a specific node or all nodes in a cluster.
Running nodetool repair -pr on a specific node repairs the subset of data owned by the target node throughout the cluster.
The subset of data is determined by the token range owned by the node.
This approach still involves all replicas, but it limits the repair to a specific token range.
Running nodetool repair on all nodes repairs the subset of data owned by any node in the cluster if the comparison process detects an inconsistency.
This approach involves all replicas and can be much more resource intensive, but it ensures that all ranges are repaired at once.
The node from which you run nodetool repair becomes the coordinator node for the repair operation.
The peer nodes are the nodes involved in the repair depending on the target of the nodetool repair command (one node’s token range or all nodes).
To build the tree, a major compaction, also known as a validation compaction, runs on the peer nodes. The validation compaction reads and generates a hash for every row in the stored tables, adds the result to a Merkle tree, and returns the tree to the coordinator node. For any given replica set, the database runs validation compaction on only one replica at a time. Other ways to refine or modify the repair scope include datacenter-wide repairs, sequential or parallel repairs, and incremental or full repairs For more information, see the following:
Hinted handoff
Occasionally, a node can become unresponsive while data is being written. Reasons for unresponsiveness include hardware problems, network issues, or overloaded nodes that experience long garbage collection (GC) pauses. If enabled, the hinted handoff feature allows the database to continue fulfilling write requests even when the cluster is operating at reduced capacity, storing undelivered writes for replay when down nodes are back online.
If the gossip-based failure detector marks a node as down and hinted_handoff_enabled is set to true in the cassandra.yaml file, the coordinator node stores missed writes as hints for a period of time.
The coordinator node can also generate hints for missed immediate writes that weren’t detected by the failure detector through gossip.
|
Hinted handoff isn’t a replacement for manual repair because hardware failure is inevitable, and hints can be lost in disaster scenarios. Depending on the severity of the outage, undelivered hints on dead nodes can be lost, and historical data across the cluster can make it difficult or impossible to reconstruct the missing data accurately. |
Hints storage
Nodes store hints in a local /hints directory with the following information for each hint:
-
Target ID for the downed node
-
Hint ID that is a time UUID for the data
-
Message ID that identifies the HCD version
-
The data itself as a blob
Hints are flushed to disk every ten seconds to minimize stale hints. When gossip discovers a node is back online, the coordinator replays each remaining hint to write the data to the newly-returned node, then deletes the hint file.
If a node is down longer than the value of the max_hint_window parameter, the coordinator stops writing new hints.
This parameter is set in cassandra.yaml, and the default is 3 hours.
When removing a node from the cluster by decommissioning the node or by using the nodetool removenode command, the database automatically removes hints targeting the node that no longer exists and removes hints for dropped tables.
Coordinator node hints
Every 10 minutes, the coordinator node checks for immediate write hints that aren’t detected through gossip.
If a replica node is overloaded or unavailable, but not yet marked as down by the failure detector, then most or all writes to that node will fail after the write_request_timeout, which is 2 seconds by default.
When this occurs, the coordinator attempts to complete the write request with other replicas, or it marks the request as failed and returns an error to the client.
If the request can be completed with other replicas, the coordinator node stores a hint for the down node.
If the request cannot be completed, either due to unavailable replicas or other reasons, the coordinator returns an error to the client, performs no writes on any replicas, and doesn’t store a hint.
For more information, see Write request coordination.
If several nodes experience brief outages simultaneously, which aren’t detected by the failure detector, substantial memory pressure can build up on the coordinator node.
The coordinator node tracks how many hints it is currently writing.
If the number of hints increases too much, the coordinator refuses writes and throws an OverloadedException error.
|
Because any node can be a coordinator node, a node that goes down might have stored undelivered hints.
Data on down nodes becomes stale after an extended outage.
If a node has been down for an extended period of time, run anti-entropy repair on the entire cluster ( |
Consistency level can prevent hint writes
The consistency level of a write request affects whether HCD writes hints and whether the write request can be fulfilled with the remaining replicas.
In general, your cluster should have enough nodes and a large enough replication factor to avoid write request failures from a small number of node outages.
As a best practice, odd-numbered replication factors and node counts help ensure that a majority of replicas are available to satisfy common consistency levels like QUORUM.
For example, at consistency level ONE, the coordinator only requires a successful acknowledgement from one replica to consider the write successful.
If the replication factor is greater than 1 but there are too many unavailable nodes, the coordinator cannot satisfy consistency level ONE.
The request fails with UnavailableException, and no hint is written because the request failed outright.
If there are sufficient replicas to meet the consistency level, the request succeeds (in the absence of other errors) and hints are written for any down replicas.
If the replication factor is 1, the coordinator node is also the node responsible for storing the data. If that node goes down, it cannot write the data or store a hint for itself. Even if there are multiple nodes in the cluster, the replication factor, by design, delivers the write to the single replica only.
To force the database to accept writes even when all replicas are down, you can use consistency level ANY.
Writes at ANY are considered durable and remain readable until an appropriate replica target becomes available to receive the hint replay.
However, if all nodes are down for too long, the hints can become stale.
At that point, extended outage repairs override any short-term repair that would be sufficiently handled by hinted handoff.
Read repair
When a read request involves multiple replicas, the coordinator node compares the data returned from each replica. If a replica’s response doesn’t match the other nodes, then the coordinator node triggers read repair to patch the mismatched replica. The coordinator node uses records with the most recent timestamp for the read repair and the response to the client.
Read repair limitations
|
Read repair isn’t the same as anti-entropy repair. Don’t rely on read repair as your sole means of synchronizing replicas. |
-
Read repair doesn’t propagate expired tombstones, and it doesn’t consider expired tombstones when selecting the most recent write for a given record. This can result in resurrected deletes (zombies) when repairing from replicas that contain missed deletes due to extended downtime. For more information, see Deletes and tombstones.
-
Read repair only includes the replicas selected for a read; all other replicas are excluded from the repair.
-
Read repair is never triggered at
ONEorLOCAL_ONEbecause those consistency levels read directly from one replica without comparing replicas. If replicas aren’t routinely repaired by other means, responses can be inconsistent when different nodes are selected for similar reads. -
Read repair runs in the foreground and blocks application operations, including responding to the original read request, until the repair is complete. If a cluster is stable and receives routine anti-entropy repairs and other consistency maintenance, read repairs shouldn’t occur frequently. If nodes frequently require read repair, there are underlying stability or consistency issues that need to be addressed.