Design bounded partition keys from real query patterns, growth distributions, skew, and bucket math while making single- vs multi-partition read costs explicit.
Partition Key Design: Distribution, Bounded Partitions, and Single-Partition Queries
Model AtlasMart documents with nested objects, arrays, dynamic fields, validation, and compatibility-safe schema evolution while distinguishing missing from explicit null.
Learning outcomes
Lesson 1 established that the partition key is the routing unit. Now AtlasMart must prevent a “correct” key from becoming operationally unbounded. A bounded partition has a defensible upper distribution for rows, bytes, tombstones, and request work under expected skew. A composite partition key uses multiple values—such as tenant, device, and day—to balance locality with growth. The design objective is not the fewest possible partitions; it is predictable partition-local work with enough distribution to survive real traffic.
Estimate partition growth from event rate, retention, and serialized row size.
Identify keys that are unbounded or hot under skew.
Use bucketing/composite keys to bound growth while preserving useful locality.
Distinguish one-partition point/range reads from multi-partition history queries.
1. Key design starts from the request and ends with a size distribution
Suppose AtlasMart telemetry writes are append-heavy. A key of
tenant_id alone answers “all telemetry for this
tenant” locally, but that is precisely the problem: the
partition grows with every device and every day. A key of
(tenant_id, device_id, day) bounds the unit by one
device-day. A 30-day device history now touches 30 predictable
partitions, which is often safer than creating one partition
whose size grows forever.
| Candidate key | Locality | Growth risk | Operational consequence |
|---|---|---|---|
| tenant_id | All tenant telemetry together | Unbounded across devices and retention | Hot/huge partition, wide failure blast radius |
| tenant_id + device_id | All history for one device | Still unbounded over time | Eventually large; compaction/repair cost grows |
| tenant_id + device_id + day | One device-day | Bounded by daily rate | More partitions for long-history reads |
| tenant_id + device_id + hour | One device-hour | Tighter bound | More fan-out for daily reads |
2. Estimate first; do not discover the partition limit during an incident
The exact “too large” threshold is product, workload, hardware, compaction, query, and version dependent. The correct practice is to model your distribution. Estimate rows per bucket, serialized bytes, tombstone density, index overhead, replicas, compaction headroom, and the high-percentile—not just average—device rate. Cassandra's current documentation notes that an exceptionally large partition can cause an SSTable to exceed its target size because that partition is not split across SSTables. That is an implementation example of the general design risk, not a universal numeric limit.
from collections import defaultdict
import hashlib
NODES = [f"n{i}" for i in range(8)]
BYTES_PER_EVENT = 420
def route(pk):
h = hashlib.blake2b(repr(pk).encode(), digest_size=8).digest()
return NODES[int.from_bytes(h, "big") % len(NODES)]
# AtlasMart telemetry planning inputs, not measured product limits.
devices = {
"normal-1": 12_000, # events/day
"normal-2": 8_000,
"celebrity-display": 480_000,
}
DAYS = 30
# Bad key: tenant only. Every device/day piles into one growing partition.
bad_rows = sum(rate * DAYS for rate in devices.values())
bad_bytes = bad_rows * BYTES_PER_EVENT
print("bad tenant-only partition rows:", bad_rows)
print("bad tenant-only estimated GiB:", round(bad_bytes / (1024**3), 3))
# Better: tenant + device + day. Growth is bounded by one device-day.
parts = {}
load = defaultdict(int)
for device, rate in devices.items():
for day in range(DAYS):
pk = ("tenant-7", device, f"2026-08-{day+1:02d}")
parts[pk] = rate
load[route(pk)] += rate
largest_pk, largest_rows = max(parts.items(), key=lambda kv: kv[1])
print("bounded partition count:", len(parts))
print("largest partition:", largest_pk, "rows:", largest_rows,
"MiB:", round(largest_rows*BYTES_PER_EVENT/(1024**2), 1))
print("one device-day query partitions touched:", 1)
print("30-day history query partitions touched:", 30)
print("rows routed per node:", dict(sorted(load.items())))
The planning model deliberately includes a “celebrity display” that emits far more events than normal devices. The bad tenant-only key produces one enormous growing partition. Daily bucketing creates 90 partitions in this 30-day/three-device scenario and keeps a normal request at one partition, while a month query intentionally pays 30 partition reads. Node totals expose distribution skew without pretending the hash model is a real Cassandra benchmark.
3. Single-partition queries are a design target, not a religious rule
Many wide-column systems are optimized for requests whose partition key is known. That gives the coordinator a deterministic routing target and bounds replica work. Multi-partition requests can still be valid when the number of buckets is small and explicit—for example, seven daily buckets for one week. The dangerous path is an unbounded or unknown fan-out where the client cannot predict how many partitions it will touch.
Track partitions/request and rows/partition as histograms. A design with median one-partition reads can still fail if one tenant produces thousands of buckets or one “celebrity” key dominates a node.
4. Wrong approach: choose the hash key only for evenness
Hashing a unique event identifier distributes writes but loses the query's routing information. The opposite mistake—using a low-cardinality tenant or region alone—preserves locality but creates hot, unbounded partitions. The repair is workload-shaped bucketing: include enough business identity to route the query and enough time/count/hash bucketing to bound the physical work. Every extra bucket must be justified by a known read fan-out and a reconstruction plan.
5. Production judgment
Before launch, record estimated p50/p95/p99 rows and bytes per partition, maximum partition request rate, bucket rollover rules, retention/TTL behavior, cross-bucket query fan-out, replica placement, and rebalance cost. Load-test skew and failure, not just uniform synthetic keys. If a single business entity can exceed the bucket limit, add another dimension—time, sub-entity, or write shard—and explicitly design the read merge. Security must remain tenant-aware across every bucket and derived table.
Check your understanding
- Why can tenant_id alone be a dangerous partition key?
- What tradeoff does daily bucketing introduce?
- Why is average partition size insufficient?
- What should a partition-size estimate include beyond row count?
- When can a multi-partition query be acceptable?
Review the answers
1. Its partition grows with every device and retention period for the tenant and can become both huge and hot.
2. It bounds each partition but long-range reads must explicitly read and merge multiple day buckets.
3. Skewed/celebrity entities can dominate tail size and load even when the mean is small.
4. Serialized bytes, tombstones, index/storage overhead, compaction/repair headroom, replication, and skew.
5. When the fan-out is small, explicit, bounded, observable, and still meets latency/failure SLOs.
Authoritative references
- Chang et al. — Bigtable: A Distributed Storage System for Structured Data (OSDI 2006) — primary research for the wide-column lineage, ordered row keys, column families, and sparse structured data.
- Apache Cassandra 5.0 — Data definition — current implementation reference for partition keys, clustering columns, composite primary keys, and clustering order.
- Apache Cassandra — Logical data modeling — query-first table design and primary-key reasoning.
- Apache Cassandra 5.0 — Tombstones — deletion markers, grace/reconciliation, repair interaction, and resurrection risk.
- Apache Cassandra 5.0 — Compaction overview — immutable SSTables, compaction, TTL/deletes, and read/space effects.
- Apache Cassandra releases — Cassandra 5.0.9 is the latest GA release as of the August 2026 review.