Every Cassandra operation that matters at scale, from where a write goes to how long a bootstrap takes and how many failures a cluster survives, is decided by a list of 64-bit numbers: the tokens each node holds. Virtual nodes (vnodes) mean each node holds many of them instead of one. That one setting changes balance, streaming, repair and availability, and it cannot be changed casually once a node has joined.
This article treats the ring as arithmetic. It shows exactly which node owns which key, gives code that computes ownership and replica sets, quantifies what many tokens buy and what they cost with a simulation whose parameters are stated, and ends with the operational procedures: allocating tokens in a new cluster, reading nodetool output and changing num_tokens without downtime. Why vnodes were introduced and how they affect repair is covered in Cassandra vnodes architecture; how a partition key becomes a token and how drivers route is in Cassandra partitioning.
The ring in four rules
- The default Murmur3Partitioner hashes each partition key to a signed 64-bit token, from -263 to 263-1. Picture that range as a circle that wraps from the maximum back to the minimum.
- Every node holds a set of tokens on that circle: one per node in the old single-token design, num_tokens of them with vnodes.
- A node's token t owns the range from the previous token on the ring, exclusive, up to t, inclusive: (previous, t]. Equivalently, a key belongs to the first node token at or after the key's token, wrapping past the maximum to the smallest token.
- That owner is the primary replica. The replication strategy then walks clockwise from it to choose the remaining replicas.
The direction matters when you read tools and documentation. A range is named by its end token, and data for a key always lives at the next node token clockwise, never the previous one.
Worked example on a toy ring
Shrink the ring to 0-99 so the numbers are readable. Node A holds tokens 10 and 50, B holds 30 and 80, C holds 65 and 95. The ranges are (95, 10] wrapping through zero, owned by A; (10, 30] by B; (30, 50] by A; (50, 65] by C; (65, 80] by B; and (80, 95] by C. Adding arc lengths, A owns 15 + 20 = 35 percent, B owns 20 + 15 = 35 percent and C owns 15 + 15 = 30 percent.
A key hashing to 42 falls in (30, 50], so A is its primary replica. With a replication factor of 2 and one rack, the strategy keeps walking: the next token is 65, held by C, so the replica set is {A, C}. A key hashing to 97 is past the last token, wraps to 10 and lands on A; the walk continues to 30, giving {A, B}. Notice that A's two ranges have different partners. With vnodes, each node shares replicas with many peers, and that single fact explains both the benefits and the costs below.
Ownership and replicas in code
The following Python reproduces the rules for one datacenter, including a simplified version of how NetworkTopologyStrategy prefers distinct racks. It is a model for reasoning and testing token layouts, not a reimplementation of Cassandra.
import bisect
SPAN = 2**64
def build_ring(nodes):
"""nodes: {name: (rack, [tokens])} -> sorted [(token, name)]"""
return sorted((t, n) for n, (_, toks) in nodes.items() for t in toks)
def ownership(ring):
"""Primary ownership: each token owns (previous token, token]."""
owned = {}
for k, (tok, name) in enumerate(ring):
width = (tok - ring[k - 1][0]) % SPAN or SPAN # k = 0 wraps
owned[name] = owned.get(name, 0) + width
return {n: w / SPAN for n, w in owned.items()}
def replicas(ring, nodes, token, rf):
"""Walk clockwise; skip a node whose rack is already used until every rack is used."""
i = bisect.bisect_left(ring, (token, ""))
all_racks = {rack for rack, _ in nodes.values()}
chosen, racks_seen, skipped = [], set(), []
for k in range(len(ring)):
name = ring[(i + k) % len(ring)][1]
if name in chosen or name in skipped:
continue
rack = nodes[name][0]
if rack in racks_seen and racks_seen != all_racks:
skipped.append(name)
continue
chosen.append(name)
racks_seen.add(rack)
if racks_seen == all_racks:
while skipped and len(chosen) < rf:
chosen.append(skipped.pop(0))
if len(chosen) >= rf:
return chosen[:rf]
return chosenTwo details are easy to get wrong. The search uses bisect_left so that a key whose token equals a node token belongs to that node, matching the inclusive end of the range. And racks change the walk: with three racks and a replication factor of 3, a node is skipped if its rack already holds a replica, so replicas land on three racks even when the next tokens clockwise belong to nodes in the same rack.
Primary versus effective ownership
The function above computes primary ownership, which is what nodetool ring hints at. What determines disk usage and load is effective ownership: the fraction of all data a node holds as any replica. In a single datacenter with a replication factor of 3 and enough nodes, effective ownership sums to 300 percent across the cluster. nodetool status only shows effective ownership when you name a keyspace, because replication is defined per keyspace.
nodetool status orders_ks # "Owns (effective)" per node, RF-adjusted
nodetool ring orders_ks # every token, its node and state
nodetool describering orders_ks # each range with its replica endpoints
nodetool getendpoints orders_ks orders 42f1c9 # replicas for one partition key
-- in cqlsh
SELECT tokens FROM system.local;
SELECT peer, tokens FROM system.peers_v2;
SELECT token(order_id), order_id FROM orders_ks.orders LIMIT 5;When effective ownership differs between nodes by more than a few percent with equal hardware, look at the token allocation, then at racks: if one rack has fewer nodes than the others, each of its nodes must hold a larger share, because every rack holds a full copy when racks equal the replication factor.
What many tokens buy: balance
With random token placement, balance is a statistics problem. To measure it, the simulation assigned random tokens to 6 nodes in one rack, computed primary ownership, and recorded the largest share over 100 random seeds. A fair share is 1/6, or 0.167.
| num_tokens (random placement) | Median largest share | Worst of 100 seeds |
|---|---|---|
| 1 | 0.390 (2.3 times fair) | 0.704 |
| 16 | 0.219 (1.3 times fair) | 0.291 |
| 256 | 0.180 (1.08 times fair) | 0.192 |
One random token per node is unusable: a typical cluster has a node carrying more than twice its share. 256 random tokens per node is well balanced, which is why it was the old default. 16 random tokens is noticeably uneven, which is why 16 is only recommended together with the token allocation algorithm that places tokens deliberately instead of randomly. These numbers are for random placement; the allocator exists to do much better than the 16-token row.
Many tokens also spread streaming. When a node joins, it takes ranges from many current owners at once rather than from one neighbour, and when a node is replaced, many peers stream to it in parallel. Cassandra bootstrap and streaming covers how those sessions run.
What many tokens cost: availability
The other side is less obvious. Consider a range with a replication factor of 3 read at QUORUM, which needs 2 of its 3 replicas. If two nodes that share any range are down at once, that range is unavailable at QUORUM. With one token per node, a node shares ranges only with its ring neighbours. With many tokens, every node shares some range with almost every other node.
The simulation counted, for 12 nodes, one datacenter, a replication factor of 3, random tokens and the simplified rack logic above, how many of the 66 possible node pairs would make some range unavailable at QUORUM if both failed.
| Layout | num_tokens 1 | num_tokens 16 | num_tokens 256 |
|---|---|---|---|
| 1 rack, 12 nodes | 24 of 66 pairs | 66 of 66 | 66 of 66 |
| 3 racks of 4 nodes | - | 48 of 66 | 48 of 66 |
With vnodes in one rack, any two simultaneous node failures make some data unavailable at QUORUM. With three racks and a replication factor of 3, the 48 pairs are exactly the pairs in different racks (3 rack pairs of 4 by 4 nodes): two failures in the same rack never break QUORUM, because every range has one replica per rack. That is the practical lesson. With vnodes, racks equal to the replication factor are what turn "any two nodes" into "any one rack", which is how you can take a whole rack down for maintenance safely.
Repair has a similar cost profile: more ranges per node mean more, smaller repair units. That trade-off is covered in Cassandra repair.
The token allocator and today's defaults
Since Cassandra 3.0, a joining node can choose its tokens with an allocation algorithm that looks at the existing ring and the replication settings, and picks tokens that minimise ownership imbalance. In 3.x it was enabled with allocate_tokens_for_keyspace, which points at a keyspace whose replication it optimises for. Cassandra 4.0 added allocate_tokens_for_local_replication_factor, which takes the replication factor directly, and the 4.0 cassandra.yaml ships with num_tokens: 16 and that option set to 3. Set it to the replication factor your keyspaces actually use in that datacenter.
The allocator needs an existing ring to optimise against, so the first nodes in a new datacenter are a special case. A common operator practice, rather than a documented requirement, is to give the first node in each rack explicit initial_token values that are evenly spaced, and let the allocator place every node after them. initial_token takes effect only on a node's first start.
def evenly_spaced(num_tokens, offset=0):
step = 2**64 // num_tokens
return [-2**63 + offset + i * step for i in range(num_tokens)]
# first node in rack 1: evenly_spaced(16); rack 2: offset step // 3; rack 3: 2 * step // 3
print(",".join(str(t) for t in evenly_spaced(16)))
Changing num_tokens on a live cluster
A node stores the tokens it chose at bootstrap in its system tables and keeps them. Changing num_tokens in cassandra.yaml on an existing node does not re-split its ranges. Replacing a dead node with replace_address_first_boot also keeps the dead node's tokens. So moving a cluster from 256 to 16 tokens means building new nodes with the new setting. The standard approach is a new datacenter:
- Start a new datacenter with num_tokens: 16 and the allocator option set, seeding the first node per rack as above.
- Switch clients to LOCAL_QUORUM or LOCAL_ONE and a datacenter-aware load balancing policy pinned to the old datacenter, so nothing reads from the empty one.
- ALTER KEYSPACE each keyspace, including system_auth and system_distributed, to replicate to the new datacenter.
- Run nodetool rebuild -- old_dc on each new node, then a repair.
- Move clients to the new datacenter, remove the old one from replication, and decommission its nodes.
This needs double the hardware for the duration, which is the real cost of choosing num_tokens badly at the start.
Failure modes
- Imbalance from mixed settings: nodes with different num_tokens own proportionally different shares. That is intended for heterogeneous hardware and a mistake otherwise.
- Forgotten cleanup: after adding nodes, existing nodes still hold data for ranges they gave away until nodetool cleanup runs on each of them.
- Uneven racks: a rack with fewer nodes overloads each of them, since every rack stores a full replica set.
- Parallel bootstraps: joining several nodes at once can break consistency guarantees and confuse the allocator. Add one node at a time.
- Hot partitions: no token layout fixes one partition key carrying most of the traffic. That is a data model problem; see Cassandra consistency levels for what QUORUM and LOCAL_QUORUM require of each replica.
What to do next
- Run nodetool status with each keyspace name and record effective ownership per node.
- Count nodes per rack and confirm racks equal the replication factor with equal node counts.
- Check num_tokens and the allocator setting in every node's cassandra.yaml.
- Load your real tokens into the simulation and count the node pairs that break QUORUM.
- For a new cluster, seed the first node per rack and use 16 tokens with the allocator.
- If you are on 256 random tokens, plan the new-datacenter migration before the cluster grows further.