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.

Intermediate120–145 minutesHot-partition diagnosis and redesign labApache Cassandra 5.0.9 · Murmur3 · 16 vnodes/nodeLast reviewed: September 2026

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.

01

Distinguish monotonic source values from low-cardinality or constant partition keys and explain how Murmur3 affects each.

02

Explain why one hot partition remains one token even after cluster expansion.

03

Use controlled reads/writes and nodetool toppartitions or equivalent metrics to identify hot keys.

04

Compare load/ownership metrics with request distribution instead of treating disk balance as traffic balance.

05

Redesign an AtlasMart hot-key table with explicit bucketing/sharding tied to query and retention requirements.

Chapter baseline reviewed 7 September 2026

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.

Execution and evidence note

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

bash / PowerShell · reproducible 3-node, 16-vnode lab
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
broken design · one global partition
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.

hot-key evidence · sample requests separately from disk ownership
# 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.

Do not publish this loop as a benchmark.

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.

controlled proof · node count changes, token(global) does not
# 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.

repair · composite partition key creates independent token positions
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

  1. Why are monotonically increasing IDs not automatically a Murmur3 hotspot?
  2. Why is a constant partition key still hot?
  3. Does RF=3 split one partition across three machines?
  4. What does nodetool status load prove about request skew?
  5. 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

cleanup · Bash; remove only course-owned resources
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
cleanup · PowerShell; remove only course-owned resources
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

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.