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.
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.
Estimate partition count from business cardinality and bucket dimensions before deployment.
Separate distribution cardinality from useful read locality; more partitions are not automatically better.
Estimate rows per partition from peak event rate, bucket duration, tenant skew and retention.
Use token samples and nodetool partition-size statistics as evidence after loading a fixture.
Diagnose low-cardinality keys and excessive sharding that trades one hotspot for uncontrolled read fan-out.
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.
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.
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
# 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"
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.
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.
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
- Why can six region values stay hot on a 100-node cluster?
- Is maximum cardinality always best?
- What does tablehistograms add beyond a spreadsheet estimate?
- Why use peak rate in row-count estimates?
- 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.
# 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
- CQL data definition — primary-key grammar, partition keys, clustering columns, and clustering order.
- CQL data manipulation — primary-key restrictions, contiguous clustering slices, token queries, and ordering behavior.
- Cassandra data-modeling introduction — partition-key and clustering-key physical meaning.
-
CREATE TABLE reference
— composite partition keys and
CLUSTERING ORDER BY. - nodetool tablehistograms — table-level percentile evidence including partition-size distribution.
- nodetool tablestats — table statistics and partition-size-related metrics.