Chapter 07 · Primary Keys, Partition Keys, Clustering Columns, and Ordering
Time Bucketing, Hash Bucketing, Synthetic Shards, and Avoiding Unbounded Partitions
Bound hot-key growth with explicit time/shard dimensions while accounting for every extra partition a reader must merge.
Learning outcomes
AtlasMart's busiest marketplace tenant can generate enough events that one “tenant” partition is unacceptable even though tenant ID has high global cardinality. The remedy is not random fragmentation. Time buckets and deterministic synthetic shards let the team bound rows and write pressure while keeping read fan-out explicit.
Choose time-bucket granularity from peak rate, retention and read windows instead of folklore.
Explain deterministic hash sharding as an explicit partition-key component.
Calculate the maximum fan-out introduced by S shards and multiple time buckets.
Distinguish deterministic routing from random shard assignment.
Recognize when bucketing reduces partition size but creates too many tiny partitions or application merges.
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. Bucketing changes the partition identity
A time bucket works only when it is part of the partition key.
For events keyed by
((tenant_id,hour_bucket,shard),event_time,event_id), every hour/shard pair is a separate partition and therefore a
separate token/replica placement. The application must derive
the same bucket on writes and reads. A retention policy spanning
24 hours with four shards has a known upper bound of 96
partitions for a full-day scan; that may be acceptable or far
too expensive depending on the API.
| Technique | What it bounds | Read cost introduced | Failure mode |
|---|---|---|---|
| hour/day bucket | time growth | one partition per bucket touched | bucket too coarse remains huge; too fine creates many tiny reads |
| deterministic hash shard | hot-key write rate per partition | up to S partitions when query spans all shards | S too large creates fan-out/tail latency |
| random shard | can spread writes | reader may not know where a row lives | uncontrolled lookup fan-out / duplicate routing ambiguity |
2. Derive shard count from a capacity hypothesis
Suppose one tenant peaks at 80,000 writes/minute and an initial test shows one partition's replica set cannot meet the required p99 at that rate. Four deterministic shards target roughly one quarter of that tenant's writes per partition if the shard input itself is well distributed. This is a hypothesis to load-test, not a universal formula. The cost is that a query requiring all events for that hour now executes four partition reads and merges ordered streams.
A common application rule is
shard = stable_hash(entity_or_event_key) mod S.
The hash algorithm and S become
schema/application contracts. Changing either changes routing
and normally requires dual-read/migration logic. Do not use a
language's process-randomized hash function if it can change
between processes or versions.
3. Lab: four bounded shard partitions
# 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.events_by_tenant_hour_shard ( tenant_id uuid, hour_bucket text, shard tinyint, event_time timestamp, event_id timeuuid, event_type text, PRIMARY KEY ((tenant_id, hour_bucket, shard), event_time, event_id)) WITH CLUSTERING ORDER BY (event_time DESC, event_id DESC) AND compaction = {'class':'UnifiedCompactionStrategy'};-- Synthetic examples: application has already derived deterministic shards 0..3.INSERT INTO atlasmart_keys.events_by_tenant_hour_shard VALUES (aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa,'2026-09-07T10',0,'2026-09-07T10:01:00Z',11111111-1111-11f1-8000-000000000001,'VIEW');INSERT INTO atlasmart_keys.events_by_tenant_hour_shard VALUES (aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa,'2026-09-07T10',1,'2026-09-07T10:02:00Z',22222222-2222-11f1-8000-000000000002,'VIEW');INSERT INTO atlasmart_keys.events_by_tenant_hour_shard VALUES (aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa,'2026-09-07T10',2,'2026-09-07T10:03:00Z',33333333-3333-11f1-8000-000000000003,'CART');INSERT INTO atlasmart_keys.events_by_tenant_hour_shard VALUES (aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa,'2026-09-07T10',3,'2026-09-07T10:04:00Z',44444444-4444-11f1-8000-000000000004,'BUY');SELECT tenant_id,hour_bucket,shard,token(tenant_id,hour_bucket,shard) AS token,event_timeFROM atlasmart_keys.events_by_tenant_hour_shard;
Expected shape: four logical partitions and normally four
different token values. Token collisions are possible in theory
but not a design mechanism. The important evidence is that
shard belongs to the partition key, so each shard
is independently placed.
SELECT event_time,event_id,event_type FROM atlasmart_keys.events_by_tenant_hour_shardWHERE tenant_id=aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa AND hour_bucket='2026-09-07T10' AND shard=0;SELECT event_time,event_id,event_type FROM atlasmart_keys.events_by_tenant_hour_shardWHERE tenant_id=aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa AND hour_bucket='2026-09-07T10' AND shard=1;SELECT event_time,event_id,event_type FROM atlasmart_keys.events_by_tenant_hour_shardWHERE tenant_id=aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa AND hour_bucket='2026-09-07T10' AND shard=2;SELECT event_time,event_id,event_type FROM atlasmart_keys.events_by_tenant_hour_shardWHERE tenant_id=aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa AND hour_bucket='2026-09-07T10' AND shard=3;
4. Broken design: random shards as invisible routing state
If AtlasMart picks a random shard for every write but stores no deterministic mapping, “fetch event X” cannot derive its partition key. The reader must search all shards or maintain a separate lookup index, turning a write shortcut into read complexity. Random sharding also makes idempotent retries dangerous if the retry selects a different shard. A deterministic shard function gives retries the same destination and makes fan-out bounded and testable.
Another mistake is treating a high shard count as free capacity. With 64 shards and a 24-hour query, a reader may need 1,536 partition requests before pagination and retries. Tail latency becomes the maximum/merge behavior of many requests, not the latency of one Cassandra read.
5. Production judgment
Choose bucket duration and shard count from workload measurements: peak writes for the hottest logical key, row size, retention, read window, concurrency, RF/CL, disk/JVM pressure, compaction and repair throughput. Test both the write benefit and the read fan-out p95/p99. Keep the mapping stable and versioned in application code. If a managed Cassandra service abstracts nodes, the same logical partition and fan-out economics still apply even if topology operations differ. Bucketing also affects backup/restore and migration: changing bucket/shard rules creates a new physical keyspace of data and usually needs dual-write/backfill/rollback planning.
Verification checklist
- Each shard value appears in the partition key and maps to an independently addressable token.
- You can enumerate the exact partitions a one-hour and one-day read must touch.
- The shard function is deterministic and stable across retries/processes.
- You have tested both write distribution and multi-shard read tail latency before increasing S.
Check your understanding
- Why must hour_bucket be inside the partition key?
- What is the main read cost of four shards?
- Why is a stable hash important for retries?
- Does 64 shards automatically improve the system?
- What must be tested after choosing S?
Review the answers
1. Only partition-key components change partition identity/token placement; a clustering-only bucket would not bound physical partition growth.
2. A query spanning the logical hour may need four partition reads plus merge/paging work.
3. The same logical write derives the same shard, supporting deterministic/idempotent routing.
4. No. It can create excessive fan-out, connection/concurrency pressure and worse tail latency.
5. Hottest-key write distribution, per-shard size, read fan-out latency, retries, compaction/repair cost, and migration/rollback behavior.
# 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
Time buckets and deterministic shards turn unbounded growth or write concentration into explicit, measurable partitions—but they also create explicit read fan-out. The final lesson makes this quantitative by estimating rows and bytes before deployment and validating the estimate with Cassandra's table statistics.
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.