Chapter 07 · Partitioning and Sharding: Range, Hash, Directory, and Consistent Hashing
Consistent Hashing, Token Rings, Virtual Nodes, and Rebalancing
Build a consistent-hash token ring, compare it with modulo hashing, add virtual nodes, and quantify how topology changes remap and stream data.
Learning outcomes
Modulo hashing distributes AtlasMart keys reasonably while the
node count is fixed, but a simple
hash(key) % N rule has an operational flaw:
changing N changes the owner for a large fraction
of keys. Consistent hashing introduces a stable hash space and
moves only ranges whose ownership changes.
Build a small token ring and route keys clockwise to token owners.
Compare remapping under naive modulo hashing with consistent hashing when a node is added.
Explain token ranges, endpoints, virtual nodes (vnodes), and why multiple tokens can improve balance.
Separate reduced remapping from the real cost of streaming/rebuilding moved data.
Reason about replica placement and topology awareness on top of a partition token map.
1. From “node count” to a stable hash space
In consistent hashing, both keys and ownership points are mapped into a fixed circular hash space. A key is owned by the first token encountered clockwise (one common convention; implementations differ). Adding a node inserts new ownership points without changing the hash of every key. Only keys in the ranges newly captured by those points move.
This solves the remapping problem, not all balancing problems. One token per physical node can still produce uneven ranges; heterogeneous node capacity and skewed keys complicate balance further.
| Term | Meaning |
|---|---|
| Token | A position/boundary in the hash space. |
| Token range | The interval whose keys map to a token owner under the routing rule. |
| Endpoint/node | A physical/logical server that owns one or more token ranges. |
| Virtual node (vnode) | One of several token positions assigned to the same physical node. |
| Token map | The current mapping from token/range ownership to endpoints. |
2. Why modulo hashing remaps so much
With hash(key) % 3, node indices are 0–2. Add a
fourth node and the divisor becomes 4; many keys now produce a
different remainder even though most existing nodes did not
fail. The network/disk cost of moving those keys can dwarf the
cost of serving normal traffic.
A consistent-hash ring decouples the hash space from the number of nodes. If node D acquires a new token between A and B, only the predecessor interval for that token changes owner. The exact moved fraction depends on token placement and key distribution.
3. Virtual nodes trade smoother balance for metadata/streaming complexity
A vnode lets one physical node own multiple smaller token ranges. When a new node joins, it can take small ranges from several existing nodes rather than one large contiguous range. This often improves incremental balance and parallelizes streaming, but more ranges also increase ownership metadata and peer relationships; products choose different defaults and token allocation algorithms.
Replica placement is a separate layer: after finding the primary partition owner, a system may place additional replicas on distinct racks/zones. Consistent hashing alone does not guarantee failure-domain diversity.
4. AtlasMart lab — add a node and count moved keys
Environment: Python 3.13.5, SHA-256-derived 8-bit teaching hash, 200 synthetic SKU keys. The tiny hash space is intentionally visual and is not a cryptographic/security recommendation or production partitioner.
import hashlib
from bisect import bisect_right
from collections import Counter
SPACE = 256
KEYS = [f"sku:{i:03d}" for i in range(200)]
def h8(text):
return hashlib.sha256(text.encode()).digest()[0]
def modulo_owner(key, nodes):
return nodes[h8(key) % len(nodes)]
def ring_owner(key, tokens):
x = h8(key)
ordered = sorted(tokens)
pos = bisect_right(ordered, x)
token = ordered[pos % len(ordered)]
return tokens[token]
nodes3 = ["A", "B", "C"]
nodes4 = ["A", "B", "C", "D"]
mod3 = {k: modulo_owner(k, nodes3) for k in KEYS}
mod4 = {k: modulo_owner(k, nodes4) for k in KEYS}
mod_moved = sum(mod3[k] != mod4[k] for k in KEYS)
# one token per node, then add D at token 96
ring3 = {0:"A", 96:"B", 192:"C"}
ring4 = {0:"A", 64:"D", 96:"B", 192:"C"}
r3 = {k:ring_owner(k, ring3) for k in KEYS}
r4 = {k:ring_owner(k, ring4) for k in KEYS}
ring_moved = sum(r3[k] != r4[k] for k in KEYS)
print("MODULO HASHING")
print("keys remapped after 3->4 nodes:", mod_moved, "/", len(KEYS))
print("\nCONSISTENT-HASH RING")
print("tokens before:", ring3)
print("tokens after :", ring4)
print("keys remapped after adding D:", ring_moved, "/", len(KEYS))
print("new owner distribution:", dict(sorted(Counter(r4.values()).items())))
# virtual-node illustration: each physical node owns multiple small ranges
vnodes = {0:"A",32:"B",64:"C",96:"A",128:"B",160:"C",192:"A",224:"B"}
print("\nVIRTUAL NODES")
print("token map:", vnodes)
print("distribution:", dict(sorted(Counter(ring_owner(k, vnodes) for k in KEYS).items())))
print("lesson: reduced remapping does not mean zero streaming; moved ranges still consume disk/network")
MODULO HASHING
keys remapped after 3->4 nodes: 147 / 200
CONSISTENT-HASH RING
tokens before: {0: 'A', 96: 'B', 192: 'C'}
tokens after : {0: 'A', 64: 'D', 96: 'B', 192: 'C'}
keys remapped after adding D: 42 / 200
new owner distribution: {'A': 56, 'B': 28, 'C': 74, 'D': 42}
VIRTUAL NODES
token map: {0: 'A', 32: 'B', 64: 'C', 96: 'A', 128: 'B', 160: 'C', 192: 'A', 224: 'B'}
distribution: {'A': 78, 'B': 82, 'C': 40}
lesson: reduced remapping does not mean zero streaming; moved ranges still consume disk/network
What to notice
The modulo rule remaps a large fraction of the 200 keys after the node count changes. The ring moves fewer keys because only the ownership interval captured by the new token changes. The vnode map then shows one physical node represented at multiple positions.
“Consistent hashing means rebalancing is cheap” is incomplete. It reduces the fraction of keys whose owner changes; every moved byte still has to be copied, verified, caught up with concurrent writes, and removed from the donor safely. Compaction/index rebuilds and replica placement can multiply that work.
Check your understanding
- Why does modulo hashing remap many keys when N changes?
- What does consistent hashing stabilize?
- Why use vnodes?
- What does a token ring not guarantee?
- What should be observed during rebalancing?
Review the answers
1. Because the divisor is part of every ownership calculation, so changing N changes many remainders.
2. The hash space and existing tokens; only ranges captured/lost by changed tokens need new ownership.
3. They divide ownership into smaller ranges so load and streaming work can be spread across more physical peers.
4. Balanced request load, failure-domain-aware replicas, zero data movement, or correctness under stale ownership metadata.
5. Bytes/ranges streamed, donor/recipient load, backlog, routing epoch convergence, replica placement, errors, and tail latency.
5. Current implementation context and production judgment
The original Dynamo paper describes consistent hashing and virtual nodes, while Apache Cassandra 5.0.9 documents tokens/vnodes and warns that token count affects both balance and streaming/availability tradeoffs. Those are implementation examples; this course does not copy their defaults as universal rules.
Choose a partitioner only after measuring key distribution and workload. Test node add/remove at realistic data volume with reserved headroom. Protect token-map/admin endpoints because a forged or stale ownership map can misroute requests. Next, the chapter asks what happens when the hash/range function is mathematically fine but the traffic itself is extremely skewed.
Authoritative references
- Dynamo: Amazon’s Highly Available Key-value Store — primary description of Dynamo consistent hashing and virtual nodes
- Apache Cassandra 5.0.9 — Dynamo architecture — current implementation terminology for tokens, endpoints, host IDs, and vnodes
- Apache Cassandra 5.0.9 — Topology changes — current implementation example of bootstrap/token allocation and streaming during node changes