Chapter 03 · Partitioners, Tokens, Vnodes, Token Rings, and Data Distribution
Hot Partitions, Skewed Keys, Monotonic Workloads, and Why More Nodes Do Not Fix Bad Keys
Create a deliberately hot AtlasMart partition, sample its request concentration, scale the cluster, and repair the schema instead of blaming the ring.
Learning outcomes
Balanced tokens do not guarantee balanced traffic. AtlasMart can have beautifully distributed vnode ownership and still melt one replica set if most requests target one partition key. This lesson separates token-space balance from workload-key balance and shows why adding nodes cannot split an existing hot partition automatically.
Distinguish monotonic source values from low-cardinality or constant partition keys and explain how Murmur3 affects each.
Explain why one hot partition remains one token even after cluster expansion.
Use controlled reads/writes and nodetool toppartitions or equivalent metrics to identify hot keys.
Compare load/ownership metrics with request distribution instead of treating disk balance as traffic balance.
Redesign an AtlasMart hot-key table with explicit bucketing/sharding tied to query and retention requirements.
Apache Cassandra 5.0.9 is the current GA 5.0 patch on the
official download page. The labs pin
cassandra:5.0.9 and explicitly set
CASSANDRA_NUM_TOKENS=16 so vnode behavior is
reproducible instead of inheriting an unnoticed
image/configuration default. Current Cassandra 5.0
documentation uses num_tokens: 16 as the modern
baseline; older Cassandra material often mentions 256 random
vnodes, so this chapter treats token count as a version- and
deployment-sensitive design choice rather than folklore.
The generation environment does not contain Docker or Cassandra, so Cassandra commands were checked against current official documentation but were not executed here. Exact token values, IP addresses, host IDs, ownership percentages, load, stream sizes, latency, and hot-partition samples must be captured on the learner's machine. Expected output is described as invariants or shapes, never presented as measured output.
1. Three different “skew” problems
Token imbalance means physical nodes own unequal portions of the hash space. Data-size skew means partitions or values differ enough that equal token ranges consume unequal bytes. request skew means some partition keys receive disproportionate reads/writes even if bytes are evenly distributed. Vnodes and token allocation primarily address the first problem; good partition-key/schema design is still required for the latter two.
A monotonic business identifier such as an increasing order
number is not automatically a Cassandra hotspot when it is
hashed by Murmur3, because adjacent source values scatter across
the token space. But a low-cardinality key—such as
status with only five values—or a constant
'global' partition creates very few partitions.
Hashing those values distributes only those few tokens.
| Input design | What Murmur3 can do | Remaining risk |
|---|---|---|
| Millions of distinct monotonic order IDs | Scatter them across token space | Individual very hot order can still be hot |
| Four regions | Map four values to four tokens | Only four partitions/replica sets carry all traffic |
| One global key | Map it deterministically to one token | One partition/replica set is the bottleneck |
| Time bucket + tenant/customer | Scatter many bounded partitions | Bucket width/cardinality still need workload sizing |
2. Build a deliberately bad global counter-like workload
docker network create atlasmart-cassandradocker volume create atlasmart-cass-1-datadocker volume create atlasmart-cass-2-datadocker volume create atlasmart-cass-3-datadocker run -d --name atlasmart-cass-1 --hostname atlasmart-cass-1 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_DC=dc1 -e CASSANDRA_RACK=rack1 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -v atlasmart-cass-1-data:/var/lib/cassandra cassandra:5.0.9# Wait until node 1 accepts CQL before starting peers.docker run -d --name atlasmart-cass-2 --hostname atlasmart-cass-2 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_SEEDS=atlasmart-cass-1 -e CASSANDRA_DC=dc1 -e CASSANDRA_RACK=rack2 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -v atlasmart-cass-2-data:/var/lib/cassandra cassandra:5.0.9docker run -d --name atlasmart-cass-3 --hostname atlasmart-cass-3 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_SEEDS=atlasmart-cass-1 -e CASSANDRA_DC=dc1 -e CASSANDRA_RACK=rack3 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -v atlasmart-cass-3-data:/var/lib/cassandra cassandra:5.0.9# Wait for all nodes to become Up/Normal.docker exec atlasmart-cass-1 nodetool statusdocker exec atlasmart-cass-1 nodetool version
docker exec atlasmart-cass-1 cqlsh -e "CREATE KEYSPACE IF NOT EXISTS atlasmart_hot WITH replication = {'class':'NetworkTopologyStrategy','dc1':3};"docker exec atlasmart-cass-1 cqlsh -e "CREATE TABLE IF NOT EXISTS atlasmart_hot.events_by_scope (scope text, event_id timeuuid, payload text, PRIMARY KEY (scope,event_id)) WITH CLUSTERING ORDER BY (event_id DESC);"# Every row below goes to one logical partition because scope='global'.docker exec atlasmart-cass-1 cqlsh -e "INSERT INTO atlasmart_hot.events_by_scope (scope,event_id,payload) VALUES ('global',now(),'event-a'); INSERT INTO atlasmart_hot.events_by_scope (scope,event_id,payload) VALUES ('global',now(),'event-b'); INSERT INTO atlasmart_hot.events_by_scope (scope,event_id,payload) VALUES ('global',now(),'event-c');"docker exec atlasmart-cass-1 cqlsh -e "SELECT token(scope) AS token_value FROM atlasmart_hot.events_by_scope WHERE scope='global' LIMIT 1; SELECT count(*) FROM atlasmart_hot.events_by_scope WHERE scope='global';"
The query returns one token for the entire
global partition. Rows can grow inside that
partition, but the partition cannot be spread among arbitrary
nodes without changing its partition key. Replication gives
multiple copies, not horizontal sharding of one partition.
3. Observe hot-partition evidence under load
nodetool toppartitions samples active partitions
for a table. Tool output and sampler names can evolve, so verify
nodetool help toppartitions on the pinned version
before automating it. Run the sampler while another terminal
generates repeated reads of the global partition.
# Terminal A: start a 20-second sampler (duration is milliseconds).docker exec atlasmart-cass-1 nodetool help toppartitionsdocker exec atlasmart-cass-1 nodetool toppartitions atlasmart_hot events_by_scope 20000# Terminal B: generate repeated requests while the sampler is active.# Bash example; use an equivalent PowerShell for-loop on Windows.for i in $(seq 1 200); do docker exec atlasmart-cass-2 cqlsh -e "SELECT event_id,payload FROM atlasmart_hot.events_by_scope WHERE scope='global' LIMIT 20;" >/dev/nulldone# Compare with cluster/load metadata; disk balance is not request balance.docker exec atlasmart-cass-1 nodetool status atlasmart_hot
Expected evidence shape: global should dominate
sampled reads because the test sends almost all traffic to that
partition. If the tiny run finishes too quickly or sampling does
not capture it, increase duration/request count on a disposable
lab. Do not fabricate a top-partition result.
docker exec cqlsh repeatedly measures process
startup and shell overhead as well as Cassandra. It is only a
traffic generator for the sampling lesson. Production
benchmarking requires a maintained driver/load generator,
concurrency, warmup, payload distribution, latency
percentiles, RF/CL, and JVM/disk/network context.
4. Add a node: ownership changes, the hot key does not
Now bootstrap node 4. The token map will change and streaming
may redistribute ranges, but
token('global') remains the same. Its natural
replica set may change depending on where new tokens are
allocated, yet every request to the same partition still targets
that one partition's replica set. More nodes improve total
cluster capacity only when the workload has enough partitions to
use them.
# Capture the hot partition token before scale-out.docker exec atlasmart-cass-1 cqlsh -e "SELECT token(scope) FROM atlasmart_hot.events_by_scope WHERE scope='global' LIMIT 1;"docker volume create atlasmart-cass-4-datadocker run -d --name atlasmart-cass-4 --hostname atlasmart-cass-4 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_SEEDS=atlasmart-cass-1 -e CASSANDRA_DC=dc1 -e CASSANDRA_RACK=rack1 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -v atlasmart-cass-4-data:/var/lib/cassandra cassandra:5.0.9# Wait until Up/Normal and bootstrap completes.docker exec atlasmart-cass-1 nodetool status atlasmart_hot# Token remains a function of the partition key, not node count.docker exec atlasmart-cass-1 cqlsh -e "SELECT token(scope) FROM atlasmart_hot.events_by_scope WHERE scope='global' LIMIT 1;"
The repaired design must change the partition key. For example, a telemetry/event workload can include tenant/store and a bounded time bucket; a global rate limiter can be sharded across deterministic buckets if the application can aggregate the result correctly. The exact bucket dimension comes from query semantics, traffic distribution, and retention—not a generic “add random suffix” recipe.
5. Redesign: spread work without losing the query contract
For AtlasMart events, suppose the real query is “latest events
for one store for one day.” A better table is keyed by
((store_id, event_day), event_id). Distinct stores
and days create many partitions while preserving a query-local
bounded slice.
docker exec atlasmart-cass-1 cqlsh -e "CREATE TABLE IF NOT EXISTS atlasmart_hot.events_by_store_day (store_id text, event_day date, event_id timeuuid, payload text, PRIMARY KEY ((store_id,event_day),event_id)) WITH CLUSTERING ORDER BY (event_id DESC);"docker exec atlasmart-cass-1 cqlsh -e "INSERT INTO atlasmart_hot.events_by_store_day (store_id,event_day,event_id,payload) VALUES ('store-101','2026-09-07',now(),'a'); INSERT INTO atlasmart_hot.events_by_store_day (store_id,event_day,event_id,payload) VALUES ('store-202','2026-09-07',now(),'b'); INSERT INTO atlasmart_hot.events_by_store_day (store_id,event_day,event_id,payload) VALUES ('store-303','2026-09-07',now(),'c');"docker exec atlasmart-cass-1 cqlsh -e "SELECT store_id,event_day,token(store_id,event_day) AS token_value FROM atlasmart_hot.events_by_store_day;"
The full-table SELECT here is acceptable only because the fixture is tiny and pedagogical; production Cassandra query design should read known partition keys rather than scan arbitrary partitions. Chapter 06 will derive query-first table designs systematically.
Verification checklist
- You can distinguish token skew, data-size skew, and request skew.
- You produced one deliberately hot partition and identified its single token.
- You used a sampler/metrics path rather than inferring hotness from ownership alone.
- Adding node 4 did not change the hot key's token.
- You redesigned the partition key from the intended query, not from a random sharding trick.
Check your understanding
- Why are monotonically increasing IDs not automatically a Murmur3 hotspot?
- Why is a constant partition key still hot?
- Does RF=3 split one partition across three machines?
- What does nodetool status load prove about request skew?
- What is the correct fix for a hot partition?
Review the answers
1. Distinct IDs are hashed into distributed tokens; source ordering is not preserved by Murmur3.
2. There is only one distinct partition-key value, therefore one token and one natural replica set.
3. No. It creates three replicas of the same partition; each replica stores the partition.
4. Very little by itself; load is disk/data state, while request skew requires request/partition activity metrics.
5. Redesign the partition key/bucketing around real query, cardinality, bounded-size, and retention requirements; adding nodes alone is insufficient.
Cleanup
docker rm -f atlasmart-cass-4 atlasmart-cass-1 atlasmart-cass-2 atlasmart-cass-3 2>/dev/null || truedocker volume rm atlasmart-cass-4-data atlasmart-cass-1-data atlasmart-cass-2-data atlasmart-cass-3-data 2>/dev/null || truedocker network rm atlasmart-cassandra 2>/dev/null || true
docker rm -f atlasmart-cass-4 atlasmart-cass-1 atlasmart-cass-2 atlasmart-cass-3 2>$nulldocker volume rm atlasmart-cass-4-data atlasmart-cass-1-data atlasmart-cass-2-data atlasmart-cass-3-data 2>$nulldocker network rm atlasmart-cassandra 2>$null
Summary and next step
This lesson’s concepts, evidence path, failure boundaries, and production judgment should now be explicit enough to verify rather than assume. Re-run the check-your-understanding prompts and preserve any lab evidence you need before changing or cleaning up the environment.
Next, continue to Inspect Token Ownership and Distribution with nodetool and System Tables.
Authoritative references
- Apache Cassandra 5.0 documentation — Official documentation entry point for the current 5.0 line.
- Apache Cassandra downloads — Official release page used to verify Cassandra 5.0.9 as the current GA patch.
- Dynamo architecture: token ring, vnodes, replication — Official explanation of consistent hashing, token ranges, vnodes, natural replicas, and ring membership.
- cassandra.yaml configuration — Official partitioner, num_tokens, token-allocation, and related configuration reference.
- Production token recommendations — Current guidance for vnode counts and token-allocation tradeoffs.
- CQL token() function — Official semantics for mapping partition-key values to partitioner tokens.
- cqlsh SHOW REPLICAS — Official cqlsh command for resolving a token to replicas for a keyspace.
- nodetool ring — Official ring command and keyspace requirement for topology-aware ownership.
- Topology changes and bootstrap streaming — Official bootstrap/token-allocation/streaming and cleanup guidance.
- nodetool netstats — Official command for observing streaming/network activity.
- Java Driver 4.19 TokenMap — Driver-side token ranges, node tokens, and replica lookup API.
- nodetool toppartitions — Official active-partition sampling tool; verify exact 5.0 syntax at runtime.