Chapter 07 · Primary Keys, Partition Keys, Clustering Columns, and Ordering

Partition-Key Cardinality, Distribution, Locality, and Bounded Partition Size

Turn “good distribution” into measurable partition count, growth bounds, locality, and hotspot evidence.

Intermediate90–120 minutesCardinality + partition-size evidence labApache Cassandra 5.0.9 · cqlsh/nodetool · UCSLast reviewed: September 2026

Learning outcomes

AtlasMart can write a syntactically valid table that is operationally disastrous. A partition key such as region may produce only a handful of partitions for millions of events, concentrating CPU, disk, compaction, repair and tail latency on a few replica sets. This lesson turns “high cardinality” and “bounded partitions” into estimates and measurements.

01

Estimate partition count from business cardinality and bucket dimensions before deployment.

02

Separate distribution cardinality from useful read locality; more partitions are not automatically better.

03

Estimate rows per partition from peak event rate, bucket duration, tenant skew and retention.

04

Use token samples and nodetool partition-size statistics as evidence after loading a fixture.

05

Diagnose low-cardinality keys and excessive sharding that trades one hotspot for uncontrolled read fan-out.

Chapter 07 lab baseline

The mandatory labs use the pinned cassandra:5.0.9 image. Java 17, cqlsh, and nodetool are the versions bundled by that image. The course topology is three disposable nodes (atlasmart-cass-1..3) in cluster atlasmart-course, datacenter dc1, racks rack1..rack3, 16 vnodes per node, replication factor (RF) 3, and LOCAL_QUORUM for the chapter's consistency-sensitive examples. Authentication, client TLS, internode TLS, and remote JMX are not enabled in this isolated learning network; do not copy that security posture to production. New Chapter 07 tables explicitly use UnifiedCompactionStrategy (UCS). Default table TTL is zero unless a lesson says otherwise; gc_grace_seconds is not changed. SAI/vector features are not used. No application driver is required for the mandatory lab; if you adapt the examples to an application, re-check your chosen driver's current compatibility and routing behavior.

Execution disclosure and resource path

These commands are documentation- and syntax-reviewed but were not executed in this generation environment. Treat shown output as an expected shape, then capture your own exact tokens, latencies, partition-size percentiles, and errors. A three-node local cluster commonly needs several GiB of RAM; if your machine cannot support it, use one disposable node and RF=1 to learn primary-key/clustering mechanics, but do not interpret that reduced topology as evidence about RF=3 availability or replica behavior.

1. Cardinality decides how many independent placements exist

A partition-key cardinality is the number of distinct partition-key values that exist during the relevant time horizon. Murmur3 hashing can distribute many distinct keys across token ranges, but it cannot manufacture diversity that the business key does not contain. If AtlasMart stores all events by region and has six regions, there are only six logical partitions no matter whether the cluster has three nodes or three hundred.

Candidate key Approx. distinct partitions Locality Main risk
region single digits very broad severe hotspots / unbounded growth
(tenant_id,day) tenants × active days good for daily tenant reads large tenants may still be hot
(tenant_id,hour,shard) tenants × hours × S bounded but S-way fan-out more application merge work

Cardinality is only half the decision. A random UUID as the partition key gives excellent distribution but destroys “all events for this tenant/hour” locality unless the application already knows every UUID. The right key balances distribution with query addressability.

2. Bound rows before you argue about bytes

Start with peak, not average, arrival rate. For a tenant producing 40 events/second at peak, a one-hour bucket yields about 144,000 rows before skew or retries. A one-day bucket yields about 3.46 million rows. Whether either is acceptable depends on row width, read slices, compaction, repair windows, disk, cache pressure and service-level objectives; there is no universal partition-size number that replaces measurement.

Useful predeployment worksheet

rows_per_partition ≈ peak_rows_per_second_for_key × bucket_seconds × skew_factor. Then separately estimate payload/key bytes and observe actual SSTable partition-size percentiles after representative loading. Replication factor multiplies cluster-wide copies; compression changes disk bytes; tombstones, indexes and metadata add costs not captured by a simple payload sum.

3. Lab: compare a pathological key with a bounded key

bash · verify or recreate the disposable course cluster
# If you already have the Chapter 01-06 lab, verify it first.docker exec atlasmart-cass-1 nodetool versiondocker exec atlasmart-cass-1 nodetool status# Standalone recreation path (skip existing resources as needed).docker network inspect atlasmart-cassandra >/dev/null 2>&1 || docker network create atlasmart-cassandradocker volume create atlasmart-cass-1-datadocker volume create atlasmart-cass-2-datadocker volume create atlasmart-cass-3-datadocker inspect atlasmart-cass-1 >/dev/null 2>&1 || docker 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 answers before starting peers.docker exec atlasmart-cass-1 nodetool statusdocker inspect atlasmart-cass-2 >/dev/null 2>&1 || docker run -d --name atlasmart-cass-2 --hostname atlasmart-cass-2 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_DC=dc1 -e CASSANDRA_RACK=rack2 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -e CASSANDRA_SEEDS=atlasmart-cass-1 -v atlasmart-cass-2-data:/var/lib/cassandra cassandra:5.0.9docker inspect atlasmart-cass-3 >/dev/null 2>&1 || docker run -d --name atlasmart-cass-3 --hostname atlasmart-cass-3 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_DC=dc1 -e CASSANDRA_RACK=rack3 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -e CASSANDRA_SEEDS=atlasmart-cass-1 -v atlasmart-cass-3-data:/var/lib/cassandra cassandra:5.0.9# Continue only when all three nodes report UN.docker exec atlasmart-cass-1 nodetool statusdocker exec atlasmart-cass-1 cqlsh -e "CREATE KEYSPACE IF NOT EXISTS atlasmart_keys WITH replication = {'class':'NetworkTopologyStrategy','dc1':3};"docker exec atlasmart-cass-1 cqlsh -e "DESCRIBE KEYSPACE atlasmart_keys"
sql · create bad and bounded fixtures
CREATE TABLE atlasmart_keys.views_by_region_bad (    region text,    observed_at timestamp,    event_id timeuuid,    tenant_id uuid,    product_id uuid,    PRIMARY KEY (region, observed_at, event_id)) WITH CLUSTERING ORDER BY (observed_at DESC, event_id DESC)  AND compaction = {'class':'UnifiedCompactionStrategy'};CREATE TABLE atlasmart_keys.views_by_tenant_day (    tenant_id uuid,    day_bucket date,    observed_at timestamp,    event_id timeuuid,    product_id uuid,    region text,    PRIMARY KEY ((tenant_id, day_bucket), observed_at, event_id)) WITH CLUSTERING ORDER BY (observed_at DESC, event_id DESC)  AND compaction = {'class':'UnifiedCompactionStrategy'};-- Small deterministic proof: two regions, four tenant/day partitions.INSERT INTO atlasmart_keys.views_by_region_bad VALUES ('eu','2026-09-07T10:00:00Z',11111111-1111-11f1-8000-000000000001,aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa,00000000-0000-0000-0000-000000000001);INSERT INTO atlasmart_keys.views_by_region_bad VALUES ('eu','2026-09-07T10:00:01Z',22222222-2222-11f1-8000-000000000002,bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb,00000000-0000-0000-0000-000000000002);INSERT INTO atlasmart_keys.views_by_region_bad VALUES ('us','2026-09-07T10:00:02Z',33333333-3333-11f1-8000-000000000003,cccccccc-cccc-cccc-cccc-cccccccccccc,00000000-0000-0000-0000-000000000003);SELECT region, token(region) AS token, observed_at FROM atlasmart_keys.views_by_region_bad;INSERT INTO atlasmart_keys.views_by_tenant_day VALUES (aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa,'2026-09-07','2026-09-07T10:00:00Z',11111111-1111-11f1-8000-000000000001,00000000-0000-0000-0000-000000000001,'eu');INSERT INTO atlasmart_keys.views_by_tenant_day VALUES (bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb,'2026-09-07','2026-09-07T10:00:01Z',22222222-2222-11f1-8000-000000000002,00000000-0000-0000-0000-000000000002,'eu');INSERT INTO atlasmart_keys.views_by_tenant_day VALUES (cccccccc-cccc-cccc-cccc-cccccccccccc,'2026-09-07','2026-09-07T10:00:02Z',33333333-3333-11f1-8000-000000000003,00000000-0000-0000-0000-000000000003,'us');SELECT tenant_id,day_bucket,token(tenant_id,day_bucket) AS tokenFROM atlasmart_keys.views_by_tenant_day;

The tiny fixture does not benchmark skew; it proves the placement model. The bad table has one partition per region. The bounded table creates independent tenant/day partition identities. Real distribution must be evaluated with representative tenant skew and load.

bash · flush and inspect size statistics
docker exec atlasmart-cass-1 nodetool flush atlasmart_keys views_by_tenant_daydocker exec atlasmart-cass-1 nodetool tablehistograms atlasmart_keys views_by_tenant_daydocker exec atlasmart-cass-1 nodetool tablestats -H atlasmart_keys.views_by_tenant_day# Optional: inspect your local 5.0 syntax before sampling hot partitions.docker exec atlasmart-cass-1 nodetool help toppartitions

tablehistograms exposes partition-size and cell-count distributions after data reaches SSTables. A three-row fixture is intentionally too small to produce capacity conclusions. The evidence becomes useful only after a production-like synthetic load with disclosed row widths, skew, warm-up and retention.

4. Broken extremes: too few partitions and too many shards

Low-cardinality keys concentrate load. The opposite mistake is to add a random shard to every write without a deterministic mapping or bounded reader plan. If each event chooses an arbitrary shard and the reader does not know which shard holds a requested entity, a single logical query can become an open-ended scatter/gather. A shard count must be a deliberate capacity knob with deterministic routing and a known maximum fan-out.

Adding nodes does not split an existing partition.

A single partition key hashes to one token and has one natural replica set. More cluster nodes can redistribute different partitions, but a hot partition remains one partition. Repair the data model by changing partition-key cardinality/bucketing and migrating data, not by expecting vnode count or node count to fragment a single key.

5. Production judgment

Evaluate cardinality using active tenants, time buckets, skew and future growth—not a static row count from a development database. Bound partitions for read latency, compaction, streaming, repair, snapshot/restore and failure recovery. Extremely tiny partitions can also increase metadata/index overhead and scatter reads, so “smaller” is not an infinite optimization. Observe p95/p99 partition sizes and request rates, hot-partition samplers, per-node disk/CPU/network, GC and read/write tail latency. Keep tenant isolation and authorization separate from partitioning: putting a tenant ID in a partition key does not itself enforce access control. Driver token awareness helps route known partition keys but cannot compensate for unknown/sharded fan-out.

Verification checklist

  • You can estimate distinct partition identities from business cardinality and buckets.
  • You observed different tokens for different tenant/day keys.
  • You captured partition-size statistics only after explaining the fixture size and SSTable state.
  • You can state why more nodes do not split one existing hot partition.
  • You can quantify the read fan-out created by any proposed shard count.

Check your understanding

  1. Why can six region values stay hot on a 100-node cluster?
  2. Is maximum cardinality always best?
  3. What does tablehistograms add beyond a spreadsheet estimate?
  4. Why use peak rate in row-count estimates?
  5. Why does adding a shard require a read plan?
Review the answers

1. Because only six partition-key values exist; Murmur3 distributes those six tokens but cannot split one partition into many independent partitions.

2. No. The key must also support bounded, addressable read locality; random per-row keys can force scatter/gather.

3. It provides observed partition-size/cell-count distributions from the loaded table, although the load must be representative.

4. Capacity and tail-latency failures are usually driven by peaks and skew, not average arrival rate.

5. Because readers may need to query and merge multiple shard partitions; the shard count sets fan-out cost.

bash · reset only Chapter 07 data
# Destructive only to this disposable chapter keyspace.docker exec atlasmart-cass-1 cqlsh -e "DROP KEYSPACE IF EXISTS atlasmart_keys;"# Keep the shared course containers for the next chapter, or remove them only if you want a full reset:# docker rm -f atlasmart-cass-1 atlasmart-cass-2 atlasmart-cass-3# docker volume rm atlasmart-cass-1-data atlasmart-cass-2-data atlasmart-cass-3-data# docker network rm atlasmart-cassandra

Summary and next bridge

Healthy partition keys create enough independent placements while preserving bounded query locality. Next we use clustering-column order and prefix rules to turn one selected partition into efficient contiguous slices.

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.