Chapter 07 · Partitioning and Sharding: Range, Hash, Directory, and Consistent Hashing
Why Partition Data: Capacity, Throughput, Parallelism, and Failure Isolation
Partition AtlasMart data deliberately by ownership, capacity, throughput, and failure isolation, while separating sharding from replication and exposing cross-partition coordination costs.
Learning outcomes
AtlasMart has outgrown the single-node mental model. Its order, session, inventory, and catalog workloads now exceed one machine's comfortable storage and request budget, but simply adding machines does nothing until the system defines ownership: which node is responsible for which keys, how clients find that owner, and what happens when an operation spans more than one owner.
Explain partitioning/sharding as a mapping from a logical key space to independent owners and distinguish it from replication.
Reason about capacity, throughput, parallelism, failure blast radius, and the coordination cost of cross-partition operations.
Trace a request through routing metadata to one shard and identify the evidence that proves which shard owns the key.
Diagnose the common mistake of treating “more shards” as automatic availability or transactional scale.
Run a deterministic AtlasMart routing lab and convert the result into production capacity and failure-domain questions.
Chapter 05 replicated copies of the same logical data and Chapter 06 asked how many copies must answer. Partitioning is different: it divides the logical key space so different owners hold different subsets. A production design usually combines both—partition first, then replicate each partition—but the mechanisms and failure consequences must be reasoned about separately.
1. Partitioning is an ownership function
A partition or shard is a subset of data assigned to an owner. The words are often used interchangeably, although products may reserve them for specific implementation layers. A partition key is the value from which routing is derived. A routing function maps that key to a partition; routing metadata then maps the partition to a node or replica set.
For a key order:123, a hash-based router might
compute hash(key) mod 4 and send the request to
shard 2. A range router might compare the order ID to configured
boundaries. A directory router may look up the key/tenant in a
metadata service. The important invariant is not the formula
itself; it is that all participants agree on the same current
ownership map.
| Concept | What it answers | What it does not guarantee |
|---|---|---|
| Partitioning | Which subset owns this key? | That the subset has redundant copies or survives a node failure. |
| Replication | How many copies of one partition exist? | That different partitions are balanced or easy to query together. |
| Routing metadata | Where should this request go now? | That a stale client has the newest map. |
| Cross-partition coordination | How do several owners participate in one operation? | A free distributed transaction or bounded tail latency. |
2. Why one node eventually stops being enough
Scale pressure can come from capacity (bytes no longer fit with safe headroom), throughput (CPU, disk, or network cannot serve enough requests), or parallelism (one serial bottleneck prevents independent work from progressing concurrently). Partitioning lets independent keys execute on different owners, but only if the workload actually has enough partitionable independence.
Suppose AtlasMart stores 260 GB of active operational state while a node can safely devote only 100 GB to that workload after indexes, write-ahead logs, compaction/maintenance space, and failure headroom. Four shards can make the dataset fit. That statement says nothing about replication factor, restore time, or whether a single celebrity product still sends 40% of all writes to one shard.
Never size partitions to the physical maximum of a node. Rebalancing, compaction, repair, temporary duplicates during movement, backup/restore, and one-node failure all consume headroom. Chapter 07 will treat headroom as part of the partitioning design, not an afterthought.
3. The routing path and observable state
A useful request trace is client → router → partition map epoch → owner → local storage/replicas → acknowledgement. The router should be able to expose which key, partition, owner, and routing-version decision it made. Operators should be able to compare that with server-side ownership metadata.
request_id=am-471
key=order:017
routing_epoch=42
partition=p3
owner_set=[node-c,node-f,node-h]
coordinator=node-c
result=200
This evidence proves how the specific request was routed under epoch 42. It does not prove the mapping is balanced, the replicas are caught up, or the business invariant is safe across other partitions.
4. Failure isolation is partial, not magical
If shard S2 fails while S0, S1, and S3 remain healthy, requests whose keys belong elsewhere may keep succeeding. That is useful blast-radius isolation. But requests for S2's keys are unavailable unless S2 is replicated and another copy can take ownership. A multi-partition request that needs S0 and S2 can fail even when most of the cluster is healthy.
The wrong design is to tell operators “the cluster is 75% healthy, therefore the application is healthy.” Availability must be measured by which keys and invariants are affected. Losing a tiny inventory partition for the most popular SKU can be worse than losing a much larger archival partition.
5. AtlasMart lab — route keys and observe the blast radius
Environment: Python 3.13.5 standard library, single process, deterministic SHA-256 routing, four logical shards, no database/cloud/container. The storage numbers are teaching inputs, not benchmark measurements.
import hashlib
from collections import Counter
SHARDS = ["S0", "S1", "S2", "S3"]
CAPACITY_PER_NODE_GB = 100
SINGLE_NODE_DATA_GB = 260
def stable_hash(text):
return int.from_bytes(hashlib.sha256(text.encode()).digest()[:8], "big")
def route(key):
return SHARDS[stable_hash(key) % len(SHARDS)]
orders = [f"order:{i:03d}" for i in range(1, 25)]
routes = [(k, route(k)) for k in orders]
counts = Counter(s for _, s in routes)
print("single-node capacity check:", SINGLE_NODE_DATA_GB, ">", CAPACITY_PER_NODE_GB,
"=> one node cannot hold the modeled dataset")
print("routing samples:", routes[:8])
print("key counts by shard:", dict(sorted(counts.items())))
print("partitioning spreads ownership; it does not create replicas")
failed = "S2"
affected = [k for k,s in routes if s == failed]
unaffected = [k for k,s in routes if s != failed]
print("\nFAILURE ISOLATION")
print("failed shard:", failed)
print("affected keys:", len(affected), affected[:6])
print("unaffected keys still routable:", len(unaffected))
query_keys = ["order:003", "order:010", "order:017", "order:024"]
query_shards = sorted({route(k) for k in query_keys})
print("\nCROSS-PARTITION REQUEST")
print("keys:", query_keys)
print("fan-out shards:", query_shards, "count=", len(query_shards))
print("lesson: a multi-key invariant may require coordination across owners")
single-node capacity check: 260 > 100 => one node cannot hold the modeled dataset
routing samples: [('order:001', 'S2'), ('order:002', 'S3'), ('order:003', 'S3'), ('order:004', 'S3'), ('order:005', 'S1'), ('order:006', 'S2'), ('order:007', 'S1'), ('order:008', 'S2')]
key counts by shard: {'S0': 2, 'S1': 10, 'S2': 4, 'S3': 8}
partitioning spreads ownership; it does not create replicas
FAILURE ISOLATION
failed shard: S2
affected keys: 4 ['order:001', 'order:006', 'order:008', 'order:011']
unaffected keys still routable: 20
CROSS-PARTITION REQUEST
keys: ['order:003', 'order:010', 'order:017', 'order:024']
fan-out shards: ['S1', 'S3'] count= 2
lesson: a multi-key invariant may require coordination across owners
What the evidence proves
- The 260 GB teaching dataset exceeds the modeled 100 GB single-node budget.
- Each key resolves to exactly one logical shard under the current hash function.
- Failure of one unreplicated shard affects only keys mapped there, while other keys remain routable.
- A four-key request can fan out to multiple owners, exposing coordination and tail-latency costs.
Deliberately wrong approach
“Shard it and we automatically get high availability” is wrong. Partitioning divides ownership; without replication, each shard can still be a single point of failure. The safe repair is to specify both the partition map and the replica placement/acknowledgement policy, then test failure of one partition owner and one cross-partition operation separately.
Check your understanding
- How is partitioning different from replication?
- Why can one failed shard matter even if most nodes are healthy?
- What evidence should a router expose?
- Why can cross-partition work be expensive?
- What capacity headroom belongs in the design?
Review the answers
1. Partitioning divides the key space among owners; replication creates additional copies of the same partition. They solve different problems and are often combined.
2. The lost shard may own keys required by a specific request or business invariant; node-count health is not application-level availability.
3. At minimum the key/partition decision, routing-map version or epoch, target owner/replica set, and request outcome.
4. It adds fan-out, more failure points, tail-latency exposure, and possibly distributed coordination for invariants or atomicity.
5. Space and throughput for normal growth plus repair, compaction, rebalance/migration, backup/restore, and degraded operation after a failure.
6. Production judgment
Partition when independent keys or aggregates can be owned separately and one node no longer satisfies capacity/throughput/SLO constraints. State the partition key, routing function, ownership metadata authority, replica placement, per-partition SLOs, cross-partition operations, and failure-domain assumptions. Observe per-shard bytes, request rate, CPU/disk/network, queue depth, tail latency, error rate, hot keys, and routing-map convergence.
Security boundaries also follow ownership: tenant identifiers in partition keys can leak through logs/metrics, and direct shard access may bypass a router's authorization policy. Keep administrative routing metadata authenticated and audited. Next, the course compares two concrete routing families—ordered ranges and hashing—using exactly the same AtlasMart keys.
Authoritative references
- MongoDB 8.3.8 — Sharding — current implementation example for shard keys, routing through mongos, and distributed chunks
- Google Bigtable — primary system paper illustrating ordered tablet/range partitioning and distributed tablet ownership
- Dynamo: Amazon’s Highly Available Key-value Store — primary paper separating partitioning/replication mechanisms in a highly available key-value system