Data distribution
In distributed databases, replication is essential for ensuring data reliability and fault tolerance.
In Apache Cassandra-based databases, data is organized into tables that use primary keys to identify unique records. When the replication factor is greater than 1, data is stored on multiple nodes that are known as replicas. If a replica is busy or fails, read/write requests can be forwarded to other replicas that store the same data.
In DataStax Enterprise (DSE), the total amount of data managed by the cluster is represented as a ring. DSE divides the ring into a number of token ranges equal to the number of nodes. Table data is partitioned across the nodes based on row keys that correspond to token ranges. When reading or writing a row, DSE uses the row key to determine the token and locate the node responsible for that token range.
Each node that joins the cluster (ring) is responsible for one or more token ranges. To join the ring, a node must be assigned a token that determines the node’s position in the ring and the range of data it is responsible for. The recommended way to manage token assignments is with Virtual nodes (vnodes). The intention is that each node is responsible for roughly an equal portion of the data.
The assigned token sets the terminus of the node’s range on the ring.
A node’s range spans from the node’s own token back to the preceding node on the ring.
For example, if one node is assigned token 5 and the next node is assigned token 10, the second node is responsible for the range from 10 back to 6.
Because the tokens represent points on a ring, the node with the highest token becomes the predecessor of the node with the lowest token.
This means that a node assigned a hypothetical 0 token is responsible for all possible tokens higher than the last assigned token (in the absence of redistribution).
If the first assigned token is 5 and the largest token is 20, the node assigned token 5 is responsible for tokens 5 back to 0 and the wrapping range from 21 to the highest possible token that could be generated by the partitioner.
Virtual nodes (vnodes)
Virtual nodes (vnodes) automatically allocate token ranges to each replica node, including initial token assignments and rebalancing token ranges when you add or remove nodes. Using vnodes simplifies partition distribution within a DataStax Enterprise (DSE) cluster. When a node joins the cluster, the database automatically assigns an even portion of data from the other nodes in the cluster. If a node fails, the database rebalances the load across the remaining nodes in the cluster. Additionally, DSE rebuilds dead nodes faster because every other node in the cluster participates in the rebuild.
In single-token architecture clusters, you must manually calculate and assign a single token to each node in a cluster every time the cluster topology changes. A node owns exactly one contiguous partition range in the ring space based on its one assigned token. In contrast, with vnodes, you set the number of tokens that each node owns, the database divides the ring based on the total number of vnodes, and then random allocates token ranges to each node. This means that each node is responsible for multiple, smaller partition ranges throughout the cluster, rather than a monolithic range. Consistent hashing distributes and redistributes data based on the desired number of vnodes without the need to recalculate tokens when the cluster topology changes.
Vnode configuration
To enable vnodes, set the num_tokens value in cassandra.yaml.
If your cluster has a mix of small and large capacity nodes (in terms of hardware), you can set different values for num_tokens to allocate specific proportions of data to each node.
Make sure the number of vnodes is appropriate for your cluster and workloads; more vnodes doesn’t always mean better performance. While vnodes provide considerable operational benefits, be aware that the number of vnodes you assign to any one node can impact cluster-wide operations. For example, when you increase the number of vnodes, you also increase the number of repairs that run during a repair cycle, which increases the duration of full cluster repairs. For most workloads, DataStax recommends 8 or 16 vnodes. In performance tests, 8 vnodes distributed token ranges between nodes with approximately 10 percent variance and minimal impact on performance.
| Replication factor | Approximate variance at 4 vnodes | Approximate variance at 8 vnodes | Approximate variance at 64 vnodes | Approximate variance at 128 vnodes |
|---|---|---|---|---|
2 |
17.5% |
12.5% |
3% |
1% |
3 |
14% |
10% |
2% |
1% |
5 |
11% |
7% |
1% |
1% |
Consistent hashing
When nodes are added or removed from a cluster that uses vnodes, token assignments are redistributed to rebalance the data amongst nodes. Consistent hashing helps distribute data across a cluster in a way that minimizes reorganization when nodes are added or removed. Consistent hashing uses each table’s partition key to determine hash values for data partitioning.
|
Consistent hashing applies to vnodes only. In single-token architectures, you must manually recalculate the token assignments for each node whenever the cluster topology changes. |
For example, if you have the following data:
| Name | Age | Car | Gender |
|---|---|---|---|
Jim |
36 |
Camaro |
M |
Carol |
37 |
BMW |
F |
Johnny |
12 |
null |
M |
Suzy |
10 |
null |
F |
The database assigns a hash value to each partition key. Hash values are generated by the partitioner that is set in the database configuration. For most use cases, the Murmur3Partitioner is appropriate.
| Partition key | Hash value |
|---|---|
Jim |
-2245462676723223822 |
Carol |
7723358927203680754 |
Johnny |
-6723372854036780875 |
Suzy |
1168604627387940318 |
Data is placed on each node according to the value of the partition key and the token range that the node is responsible for. For example, the previous hash values might be distributed to four nodes as follows:
| Node | Start range | End range | Partition key | Hash value |
|---|---|---|---|---|
A |
-9223372036854775808 |
-4611686018427387904 |
Johnny |
-6723372854036780875 |
B |
-4611686018427387903 |
-1 |
Jim |
-2245462676723223822 |
C |
0 |
4611686018427387903 |
Suzy |
1168604627387940318 |
D |
4611686018427387904 |
9223372036854775807 |
Carol |
7723358927203680754 |
Partitioners
The partitioner function derives a token from the partition key, typically through hashing. The cluster distributes each row of data based on the token value. This distribution works even if the tables use different partition keys, such as user names or timestamps. Because each part of the hash range receives an equal number of rows on average, the read and write requests to the cluster are evenly distributed and load balancing is simplified.
|
New nodes must use the same partitioner as the existing nodes in the cluster. If you change the partitioner for an existing cluster, you must reload all data (for example, from a snapshot). Otherwise, all existing data becomes inaccessible and effectively lost. |
Set the partitioner with the partitioner parameter in cassandra.yaml.
The following partitioners are available:
- Murmur3Partitioner (default, recommended)
-
DataStax recommends this partitioner for new clusters, and it is required to enable vnodes with replication-balanced token allocation. Legacy partitioners are included for backwards compatibility only.
This partitioner uniformly distributes data across the cluster based on MurmurHash hash values. This hashing function creates a 64-bit hash value of the partition key with a possible range from -263 to +263-1. When using this partitioner, you can page through all rows using the
TOKEN()function in a CQL query. - Legacy partitioners
-
RandomPartitioner, ByteOrderedPartitioner, and OrderPreservingPartitioner are legacy partitioners included for backwards compatibility only. Use the Murmur3Partitioner for new clusters.
The RandomPartitioner uniformly distributes data evenly across the nodes using an MD5 hash value of the row key. The possible range of hash values is from 0 to 2127-1. Because it uses a cryptographic hash, which isn’t required by the database, it takes longer to generate the hash value than the Murmur3Partitioner. When using this partitioner, you can page through all rows using the
TOKEN()function in a CQL query.The ByteOrderedPartitioner orders rows lexically by key bytes. It requires significant administrative overhead to load balance the cluster, sequential writes can cause hot spots, and balancing for one table can result in uneven distribution for another table in the same cluster.
- Bring-your-own partitioner
-
You can use any IPartitioner, including your own, as long as it is in the classpath. However, DataStax recommends the default Murmur3Partitioner.