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.

Beginner → Advanced90–110 minuteshot-key skew simulationVendor-neutral · Python 3.13.5 simulatorFree/local · no database or cloud requiredLast reviewed: August 2026

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.

01

Identify hot keys and hot partitions from a skewed request distribution instead of relying on average cluster load.

02

Explain why adding nodes cannot split one indivisible hot key automatically.

03

Compare salting/write sharding, time bucketing, caching, fan-out, and business decomposition with their read/consistency costs.

04

Recognize monotonic identifiers and celebrity objects as different paths to write concentration.

05

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.

python · hot-key and write-sharding simulator
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")
text · verified output from the simulator
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.

Do not blindly salt identifiers

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

  1. What is the difference between a hot key and a hot partition?
  2. Why does adding nodes not split one hot key?
  3. What does salting change?
  4. Why can caching fail to fix a write hotspot?
  5. 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

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.