Chapter 07 · Partitioning and Sharding: Range, Hash, Directory, and Consistent Hashing
Hot Partitions, Hot Keys, Monotonic Identifiers, Celebrity Data, and Write Concentration
Detect hot keys and partitions under skewed traffic, then compare salting, bucketing, caching, fan-out, and decomposition with their read and consistency costs.
Learning outcomes
AtlasMart adds nodes, but a livestream promotion sends hundreds of writes per second to one product record. The cluster has plenty of aggregate CPU and disk, yet one partition queue saturates. This is the difference between cluster capacity and per-key/per-partition capacity.
Identify hot keys and hot partitions from a skewed request distribution instead of relying on average cluster load.
Explain why adding nodes cannot split one indivisible hot key automatically.
Compare salting/write sharding, time bucketing, caching, fan-out, and business decomposition with their read/consistency costs.
Recognize monotonic identifiers and celebrity objects as different paths to write concentration.
Run a deterministic skew lab and verify both the hotspot and the cost of a repair strategy.
1. Averages hide skew
A hot key receives disproportionate traffic relative to other keys. A hot partition is an owner whose total workload saturates a limiting resource. One hot key can cause a hot partition, but a partition can also become hot because many medium-hot keys cluster there or because a range receives all new monotonic IDs.
Zipf-like popularity—where a few items receive most requests—is common in catalogs, social feeds, rate-limit counters, counters for viral content, and tenancy where one customer is much larger than others. Capacity planning must therefore include percentiles of key popularity, not just average bytes per shard.
2. Why adding nodes often fails to cure one hot key
If the routing unit is exactly product:42, every
request for that product resolves to one partition (plus its
replicas). Doubling the number of nodes can move the key to a
different owner, but does not divide that key's write
serialization or queue. The new owner becomes hot instead.
| Technique | What it can relieve | New cost / caveat |
|---|---|---|
| Salting / write sharding | Splits one logical key into several physical keys | Reads/aggregations fan in; global constraints need coordination or merge logic. |
| Time buckets | Bounds an append-heavy key by time window | Queries across time touch multiple buckets; late events complicate placement. |
| Cache | Reduces repeated read load | Does not reduce write concentration; introduces freshness/invalidation semantics. |
| Fan-out on write/read | Moves computation between write and read paths | Amplifies storage/write work or read fan-out. |
| Business decomposition | Splits ownership by a real domain boundary | May require product/domain redesign and new invariant boundaries. |
3. Monotonic IDs create a range hotspot; celebrities create a key hotspot
These failures look similar in dashboards but have different mechanisms. A monotonically increasing timestamp/order ID overloads the newest range; automated range splits or hashing can spread future writes. A celebrity key remains one logical value even under perfect hash distribution; only changing the logical routing unit or operation semantics can distribute its writes.
Therefore diagnosis should record top keys, key-to-partition mapping, per-partition request rate, queue depth, CPU/disk/network, and tail latency. A CPU graph by node alone is not enough.
4. AtlasMart lab — deterministic skew and a salted repair
Environment: Python 3.13.5, 1,000 deterministic synthetic requests, SHA-256 routing. The traffic mix is a teaching workload, not a measured production distribution.
import hashlib
from collections import Counter
NODES = ["N0", "N1", "N2", "N3"]
def route(key, nodes=NODES):
h = int.from_bytes(hashlib.sha256(key.encode()).digest()[:8], "big")
return nodes[h % len(nodes)]
traffic = []
traffic += ["celebrity:42"] * 420
traffic += ["flash-sale:sku-9"] * 160
for i in range(420):
traffic.append(f"user:{i:04d}")
key_counts = Counter(traffic)
node_counts = Counter(route(k) for k in traffic)
print("TOP KEYS:", key_counts.most_common(5))
print("NODE LOAD, 4 NODES:", dict(sorted(node_counts.items())))
print("celebrity key owner:", route("celebrity:42"), "requests=", key_counts["celebrity:42"])
nodes8 = [f"N{i}" for i in range(8)]
node_counts8 = Counter(route(k, nodes8) for k in traffic)
print("\nWRONG FIX: ADD NODES")
print("NODE LOAD, 8 NODES:", dict(sorted(node_counts8.items())))
print("celebrity key is still one indivisible routing unit on:", route("celebrity:42", nodes8))
print("\nWRITE-SHARD THE HOT KEY INTO 8 BUCKETS")
sharded = []
for i in range(420):
sharded.append(f"celebrity:42#bucket:{i % 8}")
sharded += [k for k in traffic if k != "celebrity:42"]
sharded_counts = Counter(route(k, nodes8) for k in sharded)
print("NODE LOAD AFTER SALTING:", dict(sorted(sharded_counts.items())))
print("read fan-in buckets for celebrity aggregate: 8")
print("tradeoff: write concentration falls, but reads/aggregation and consistency become more complex")
TOP KEYS: [('celebrity:42', 420), ('flash-sale:sku-9', 160), ('user:0000', 1), ('user:0001', 1), ('user:0002', 1)]
NODE LOAD, 4 NODES: {'N0': 90, 'N1': 549, 'N2': 262, 'N3': 99}
celebrity key owner: N1 requests= 420
WRONG FIX: ADD NODES
NODE LOAD, 8 NODES: {'N0': 39, 'N1': 62, 'N2': 214, 'N3': 45, 'N4': 51, 'N5': 487, 'N6': 48, 'N7': 54}
celebrity key is still one indivisible routing unit on: N5
WRITE-SHARD THE HOT KEY INTO 8 BUCKETS
NODE LOAD AFTER SALTING: {'N0': 144, 'N1': 62, 'N2': 214, 'N3': 150, 'N4': 103, 'N5': 120, 'N6': 101, 'N7': 106}
read fan-in buckets for celebrity aggregate: 8
tradeoff: write concentration falls, but reads/aggregation and consistency become more complex
Failure and repair
The celebrity key accounts for 420 requests and lands on one owner. Moving from four nodes to eight does not subdivide it. The repaired model gives the celebrity aggregate eight physical bucket keys, spreading its writes more broadly. But a complete celebrity-count read must now query/merge eight buckets, and any “exact global count” invariant must define merge/concurrency semantics.
Random prefixes can destroy useful locality and make operational lookup/debugging harder. Salt only when you have a specific measured hotspot and a defined reconstruction/query path.
Check your understanding
- What is the difference between a hot key and a hot partition?
- Why does adding nodes not split one hot key?
- What does salting change?
- Why can caching fail to fix a write hotspot?
- What metrics reveal skew?
Review the answers
1. A hot key is one disproportionately popular logical key; a hot partition is an owner whose aggregate workload is saturated. A hot key can cause a hot partition.
2. The routing unit remains indivisible; the key simply hashes/ranges to one owner in the new topology.
3. It converts one logical item into multiple physical routing keys, spreading writes but requiring fan-in/merge on reads.
4. A cache can absorb repeated reads, but writes still target the authoritative key/partition unless the write model changes.
5. Top-key frequency, per-partition request rate/bytes, queue depth, tail latency, CPU/disk/network, and key-to-partition ownership.
5. Production judgment
Start with instrumentation and workload evidence. If the hottest key already exceeds one partition's sustainable write budget, no amount of average cluster headroom solves it. Choose the mitigation based on invariant strength: approximate counters can shard/merge easily; inventory decrements that must never go negative may need single-owner serialization or reservation semantics rather than arbitrary sharding.
Test skew deliberately in load tests—uniform random traffic is often dangerously optimistic. Protect hot-key diagnostic data because it can reveal customer/product popularity. Next, the chapter tackles online movement itself: splitting and resharding while reads and writes continue.
Authoritative references
- MongoDB 8.3.8 — Choose a Shard Key — current implementation guidance on frequency, cardinality, monotonic growth, and query patterns
- Dynamo: Amazon’s Highly Available Key-value Store — primary partitioning/load-distribution background
- Apache Cassandra 5.0.9 — Production recommendations — implementation example of vnode/load-distribution tradeoffs