Skip to main content
Create your own

Consistent Hashing and Virtual Nodes in Cassandra

Hello! Welcome to the first lesson in our module on Cassandra.

Given your extensive background in building high-load distributed systems, you're likely familiar with the fundamental challenge of partitioning data across a cluster. This lesson delves into Cassandra's specific approach to solving this problem, which is crucial for its scalability and resilience.

Today, we will focus on the following learning outcome: Explain consistent hashing and virtual nodes in Cassandra and analyze their impact on data distribution.

We will cover:

  • How Cassandra uses a consistent hashing algorithm to map data to nodes on a token ring.
  • The evolution from a single-token-per-node architecture to the more flexible virtual node (vnode) system.
  • The critical trade-offs involved in configuring vnodes, analyzing their impact on cluster balancing, scaling, and operational performance.

This foundation is essential for understanding nearly every other aspect of Cassandra's architecture, from replication to read/write paths.

1. The Challenge of Partitioning and the Consistent Hashing Solution

In any distributed database, you need a deterministic way to assign a piece of data to a specific node. A naive approach, like using a modulo operator (hash(key) % number_of_nodes), creates significant problems. When you add or remove a node, the divisor changes, forcing a massive reshuffling of data across the entire cluster. This is untenable for large, dynamic systems.

Cassandra, inspired by Amazon's Dynamo, solves this with consistent hashing. The core idea is to map both nodes and data partitions onto the same logical space—a continuous, circular ring of hash values called the token ring.

To understand this concept in detail, please read the following introduction to Cassandra's partitioning mechanism.

The Impacts of Changing the Number of VNodes in Apache ...

This first resource, from The Last Pickle blog, provides a clear explanation of how Cassandra uses a partitioner and a token ring to assign data responsibility to nodes.

Please read the first two paragraphs of the article. Focus on understanding the concepts of the partitioner, token, and token ring, and how a node's token range is defined.

As the article explains, the partitioner (by default, Murmur3Partitioner) hashes a row's partition key to generate a token. This token's position on the ring determines which node is responsible for that data. A node "owns" the range of tokens from the previous node's token on the ring up to its own assigned token.

The key benefit of this approach is that when a node is added or removed, it only affects its immediate neighbors on the ring. Data only needs to be streamed to or from adjacent nodes, rather than being reshuffled across the entire cluster.

2. The Original Architecture: Single Token Per Node

In early versions of Cassandra, the implementation of this model was straightforward: each physical node in the cluster was assigned exactly one token. While this can lead to a perfectly balanced cluster if tokens are calculated and assigned correctly, it presents significant operational challenges.

The next section of the same article discusses this pre-vnode era.

The Impacts of Changing the Number of VNodes in Apache ...

This section describes the historical context of single-token architecture, highlighting the difficulties associated with scaling and rebalancing.

Please read the section titled 'Back in the day…'. Pay attention to the manual process of token calculation and the operational pain of growing a cluster using nodetool move.

As you read, the main takeaway is the operational inflexibility. Adding a single node to an existing cluster would unbalance the token distribution, creating "hot spots" where the new node and its neighbor have smaller ranges than others. Rebalancing required manually moving tokens, an expensive and tedious streaming operation. This made organic, incremental scaling very difficult, a major drawback for the kind of high-growth systems you've managed.

3. Vnodes: A More Flexible Approach

To address the rigidity of the single-token model, Cassandra 1.2 introduced virtual nodes (vnodes). Instead of a physical node owning one large, contiguous token range, it now owns many smaller, non-contiguous ranges.

This seemingly simple change has profound implications for cluster management and data distribution.

The Impacts of Changing the Number of VNodes in Apache ...

The following two resources explain what vnodes are and their benefits. The first introduces the concept and its goals, while the second provides a structured list of advantages and disadvantages.

First, read the section 'vnodes to the rescue' in The Last Pickle article. This covers the motivation and the introduction of the num_tokens setting.

Dynamo | Apache Cassandra Documentation

Now, read this section from the official Cassandra documentation for a concise summary of vnodes and their pros and cons.

Read the section 'Multiple Tokens per Physical Node (vnodes)'. Focus on the list of benefits and disadvantages.

By breaking a node's ownership into many small virtual nodes, Cassandra achieves several key goals:

  • Automated and Smoother Scaling: When a new node joins, it automatically gets assigned its num_tokens vnodes. It "steals" many small token ranges from nodes all across the ring, leading to a more balanced cluster without manual intervention.
  • Faster Rebuilds and Better Failure Isolation: If a node fails, the responsibility for its vnodes is distributed among many remaining nodes in the cluster, not just its two immediate neighbors. This spreads the extra load more evenly and speeds up recovery.

However, vnodes are not a silver bullet. The number of vnodes, configured via num_tokens in cassandra.yaml, introduces a critical tuning trade-off.

The num_tokens Trade-off

The default value for num_tokens was 256 for a long time. While this random token allocation simplified scaling, it came with significant downsides that are crucial to understand when designing a production system.

The Impacts of Changing the Number of VNodes in Apache ...

This next reading explores the practical consequences and 'fine print' of using vnodes, particularly the issues with both very high and very low num_tokens values.

Please read the sections 'Remember to read the fine print', 'Pics or it didn’t happen', and 'Too many vnodes spoil the cluster'. Focus on the analysis of how num_tokens affects token range balance, availability during outages, and the performance of streaming operations like repair.

Let's summarize the key trade-offs you just read about:

  1. Too few vnodes (e.g., < 32) with random allocation: The law of large numbers doesn't apply, and you can end up with significant data imbalance. One node might randomly be assigned vnodes that correspond to much larger total token ranges than another, creating a hot spot.
  2. Too many vnodes (e.g., 256):
    • Increased operational overhead: Cluster-wide operations like repair, bootstrap, or decommissioning become slower. Cassandra initiates a separate streaming session for each token range. With 256 vnodes per node, this creates a massive number of small, sequential tasks, and the overhead adds up, especially on clusters with terabytes of data per node.
    • Reduced availability in certain failure scenarios: While a single node failure is handled well, a multi-node failure has a higher probability of taking out all replicas for a specific, small token range. The more you slice the ring, the more combinations of node failures can lead to localized data unavailability.
    • Secondary index performance degradation: Queries on secondary indexes are fanned out to all nodes. Each node must then scan the data for all of its assigned token ranges. More vnodes mean more individual ranges to check, increasing latency.

Your experience with low-latency and high-load systems makes the impact of slow repair operations or unpredictable availability particularly relevant. The convenience of easy scaling with a high num_tokens count comes at a direct cost to operational performance and resilience.

4. Modern Vnodes: Intelligent Token Allocation

Recognizing these drawbacks, the Cassandra community developed more sophisticated token allocation strategies. Modern versions of Cassandra no longer rely on pure randomness.

The Impacts of Changing the Number of VNodes in Apache ...

This final reading covers the improvements made in Cassandra 3.0 and 4.0 that allow for the benefits of vnodes without the severe performance penalties.

Please read the section 'A new hope'. Focus on the concept of the replica-aware token allocation algorithm and the new default settings in Cassandra 4.0.

The key takeaway here is the shift to replica-aware token allocation. Instead of randomly scattering vnodes, Cassandra can now intelligently calculate token assignments to ensure an even distribution of data, even with a much smaller num_tokens value.

As of Cassandra 4.0, the default num_tokens is 16, and a new intelligent allocation algorithm is enabled by default. This provides a much better out-of-the-box experience, combining the operational ease of vnodes with the balance and performance of a low token count. It's a prime example of how the practical implementation of a theoretical concept evolves to address real-world operational challenges.

Conclusion

In this lesson, we've dissected how Cassandra partitions data to achieve horizontal scale.

Key Takeaways:

  • Consistent Hashing: Cassandra avoids the pitfalls of naive hashing by mapping nodes and data to a logical token ring. This minimizes data movement when the cluster size changes.
  • Single-Token Architecture: The original model was simple but operationally rigid, making scaling and rebalancing difficult and expensive.
  • Virtual Nodes (Vnodes): Assigning multiple token ranges per physical node automates scaling and improves load distribution during failures.
  • The num_tokens Trade-off: The number of vnodes is a critical parameter. A high count simplifies scaling but hurts operational performance (e.g., repairs) and can impact availability. A low count, with older random allocation, leads to data imbalance (hot spots).
  • Modern Approach: Current Cassandra versions use intelligent, replica-aware token allocation, allowing for a low num_tokens value (e.g., 16) that achieves good balance without the performance overhead of the older, high-vnode approach.

Understanding this mechanism is fundamental to designing, operating, and troubleshooting a Cassandra cluster effectively. The choice of num_tokens (though now better defaulted) and awareness of the underlying allocation strategy directly impact the performance and stability of the high-load systems you aim to build.

Preview of the Next Lesson:

Now that we know how data is partitioned across nodes, the next logical question is: how is data replicated to ensure durability and availability? In the next lesson, we will cover Cassandra's replication strategies and replication factor, including how to configure them across multiple datacenters.

Can't find a good explanation? Sign up and we'll make it for you

Sign up