Chapter 07 · Partitioning and Sharding: Range, Hash, Directory, and Consistent Hashing
Resharding, Split/Merge Operations, Online Rebalancing, and Capacity Headroom
Trace online split/reshard state through copy, catch-up, validation, routing-epoch cutover, rollback, and capacity headroom under failure.
Learning outcomes
AtlasMart has identified an overloaded range and wants to move half of it from shard S0 to S1 without stopping checkout. The operation is not “copy files and change a config.” While the copy runs, clients keep writing. Routers may cache old ownership. The recipient must catch up, ownership must switch exactly once under a newer epoch, and rollback must remain possible.
Trace an online split/move through planning, snapshot copy, change catch-up, validation, routing-epoch cutover, and cleanup.
Explain dual-ownership/forwarding windows and why stale routing metadata must be detected or redirected.
Quantify why migration requires spare disk/network/CPU and degraded-mode headroom.
Diagnose unsafe early cutover that routes reads to a recipient before data is caught up.
Design throttling, verification, rollback, and observability gates for online resharding.
1. Online movement is a state machine
A safe reshard/move usually has phases even if a product gives them different names: plan → snapshot/clone → capture concurrent changes → catch up → validate → publish new ownership epoch → drain old routes → cleanup donor. Some systems forward writes, some replay a log/oplog, some temporarily duplicate writes, and some block a short final window. The mechanism matters more than the command name.
The routing map should be versioned. An epoch here is a monotonically increasing configuration version. Requests carrying epoch 17 after ownership has moved to epoch 18 should not be silently accepted by two independent owners. They should be redirected, forwarded, fenced, or otherwise reconciled according to the system's protocol.
2. Wrong cutover: metadata changes before state is ready
The simplest failure is publishing “keys 500–999 now belong to S1” before S1 contains the snapshot and concurrent changes. A client with fresh metadata correctly routes key 700 to S1 and gets “not found,” even though the donor still has the row. This is a routing-consistency failure caused by an unsafe move sequence.
During migration, two nodes may physically hold copies of the moving range. The protocol still needs one authoritative write rule (or a strictly defined forwarding/log-catch-up mechanism). Two independently writable owners without fencing/version semantics create lost updates or split-brain state.
3. Capacity headroom must cover the failure case
Rebalancing duplicates data temporarily and consumes network, disk I/O, CPU, caches, and write-ahead/change-log space. Plan for the worst reasonable topology, not just steady state. Three 100-unit nodes at 70% aggregate utilization look comfortable while healthy, but if one node fails, the surviving two would need to absorb 210 units—105 each—before accounting for migration amplification.
| Resource | Why movement increases it | Operational signal |
|---|---|---|
| Disk capacity | Donor + recipient overlap; indexes/compaction may duplicate space | free bytes, temp bytes, compaction backlog |
| Network | Snapshot/stream + change catch-up competes with client traffic | stream Mbps, retransmits, client p99 latency |
| CPU/I/O | Checksums, index builds, compaction, serialization | CPU steal/load, disk queue/latency |
| Change log | Concurrent writes accumulate until recipient catches up | oplog/WAL/change backlog age and bytes |
| Routing metadata | Old and new maps coexist across clients/caches | redirects, stale-epoch rejects, map convergence |
4. AtlasMart lab — unsafe vs gated split
Environment: Python 3.13.5 deterministic in-memory simulation. It models range ownership and catch-up semantics; it does not reproduce a specific database's resharding implementation.
from dataclasses import dataclass
@dataclass
class Range:
lo: int
hi: int
owner: str
epoch: int
rows = {k: f"v{k}" for k in range(0, 1000, 100)}
old = Range(0, 1000, "S0", 17)
left = Range(0, 500, "S0", 18)
right = Range(500, 1000, "S1", 18)
print("START routing:", old)
print("rows on donor:", sorted(rows))
# WRONG: publish new routing before the recipient has a copy.
recipient = {}
print("\nWRONG CUTOVER")
print("publish epoch 18 immediately; key 700 routes to S1")
print("recipient has key 700?", 700 in recipient)
print("result: a routed read can miss data that still exists on donor")
print("\nSAFE ONLINE MOVE")
print("1 freeze move plan at source epoch 17")
# snapshot copy right half
recipient = {k:v for k,v in rows.items() if k >= 500}
print("2 snapshot cloned to S1:", sorted(recipient))
# concurrent writes are logged/forwarded while copy happens
change_log = [(700, "v700-new"), (900, "v900-new")]
for k,v in change_log:
rows[k] = v
recipient[k] = v
print("3 catch-up changes applied:", change_log)
# validate equality for moving range
source_right = {k:v for k,v in rows.items() if k >= 500}
print("4 validation equal?", source_right == recipient)
print("5 publish epoch 18 only after catch-up/validation")
print("new ranges:", left, right)
print("6 keep rollback metadata until old clients/routing caches converge")
print("\nCAPACITY HEADROOM")
cluster_capacity = 300
used = 210
remaining_nodes_after_failure = 2
required_per_remaining = used / remaining_nodes_after_failure
print("cluster used before failure:", used, "/", cluster_capacity, "=", used/cluster_capacity)
print("after one of three equal nodes fails, surviving nodes would each need:", required_per_remaining,
"units against 100-unit capacity")
print("safe steady-state utilization must leave room for failure plus migration amplification; 70% here is not enough")
START routing: Range(lo=0, hi=1000, owner='S0', epoch=17)
rows on donor: [0, 100, 200, 300, 400, 500, 600, 700, 800, 900]
WRONG CUTOVER
publish epoch 18 immediately; key 700 routes to S1
recipient has key 700? False
result: a routed read can miss data that still exists on donor
SAFE ONLINE MOVE
1 freeze move plan at source epoch 17
2 snapshot cloned to S1: [500, 600, 700, 800, 900]
3 catch-up changes applied: [(700, 'v700-new'), (900, 'v900-new')]
4 validation equal? True
5 publish epoch 18 only after catch-up/validation
new ranges: Range(lo=0, hi=500, owner='S0', epoch=18) Range(lo=500, hi=1000, owner='S1', epoch=18)
6 keep rollback metadata until old clients/routing caches converge
CAPACITY HEADROOM
cluster used before failure: 210 / 300 = 0.7
after one of three equal nodes fails, surviving nodes would each need: 105.0 units against 100-unit capacity
safe steady-state utilization must leave room for failure plus migration amplification; 70% here is not enough
Verification gates
- Recipient snapshot contains every row in the moving range at the chosen source point.
- Concurrent changes are replayed/forwarded until lag reaches the system's cutover requirement.
- Source and recipient checksums/counts/version summaries match for the moving range.
- Only then is routing epoch 18 published.
- Old-epoch traffic is redirected/fenced while caches converge.
- Donor cleanup waits until rollback/read verification windows close.
Deliberately wrong approach
“Change the router first so new traffic helps warm the recipient” creates a client-visible hole if the recipient is not complete. The lab shows key 700 routed to an empty recipient. The repair is a gated state machine with snapshot/catch-up/validation before cutover.
Check your understanding
- Why is a routing epoch useful?
- What is the purpose of a catch-up phase?
- Why keep rollback metadata after cutover?
- Why is 70% steady-state utilization potentially unsafe in a three-node cluster?
- What should throttle a rebalance?
Review the answers
1. It lets nodes/clients detect stale ownership decisions and prevents silent acceptance under conflicting maps.
2. It applies writes that occurred after the initial snapshot so the recipient reaches the cutover state.
3. Old clients and hidden correctness defects may appear after the map changes; rollback requires knowing the prior owner/state until confidence is established.
4. After one node fails, remaining nodes may need more than their capacity before migration overhead is counted.
5. Client SLO impact and resource saturation—disk/network/CPU, change backlog, p99 latency/errors—not an arbitrary fixed transfer rate alone.
5. Current implementation context
MongoDB 8.3.8's current stable-series resharding documentation is a concrete example of the resource reality: it warns about storage, I/O, CPU, and oplog pressure and describes donor/recipient/coordinator roles. Apache Cassandra 5.0.9's current topology-change documentation similarly treats node bootstrap as token-range streaming. These product workflows differ, but both reinforce the vendor-neutral lesson that online ownership changes consume real capacity and need coordination.
6. Chapter synthesis and production judgment
Chapter 07 started with “which owner holds this key?” and ended with “how does ownership change safely while traffic continues?” The complete partitioning decision now includes key shape, locality, skew, routing metadata, replica placement, hot-key mitigation, growth model, movement protocol, failure blast radius, and headroom.
Production runbooks should define preconditions (free disk, replica health, repair/backlog state), canary range, transfer throttles, verification checks, routing-epoch monitoring, abort conditions, rollback, and post-move cleanup. Never perform failure injection or resharding experiments against production data merely to learn the mechanism; use isolated environments or formally reviewed game days. Chapter 08 next explains how nodes learn membership/topology changes, suspect failures, gossip metadata, and avoid two active owners.
Authoritative references
- MongoDB 8.3.8 — Reshard a Collection — current implementation example with donor/recipient roles, resource requirements, monitoring, and cutover behavior
- MongoDB 8.3.8 — Sharding — current implementation context for routing and data movement across shards
- Apache Cassandra 5.0.9 — Topology changes — current implementation example of bootstrap/token allocation and streaming
- Dynamo: Amazon’s Highly Available Key-value Store — primary background for partition ownership and incremental growth