Why this matters
- Naive modulo hashing remaps almost every key when N changes — unacceptable during cluster resize.
- Consistent hashing is standard for caches, shard routers, and load balancers.
- Virtual nodes prevent uneven ring segments when physical machines differ in capacity.
- ShardPay's shard router uses consistent hashing with 150 vnodes per physical shard.
Keys hash onto the ring; walk clockwise to find owner
Ring mechanics
ShardPay places 150 virtual nodes (vnodes) per physical shard on a hash ring from 0 to 2^64−1. An accountId hashes to a point on the ring; the router walks clockwise to the first vnode ≥ that hash — that vnode's physical shard owns the account. Adding shard 257 inserts 150 new vnodes; only keys between adjacent new and old vnodes migrate — not the entire keyspace.
Ring concepts
- Hash ring — Logical circle where both keys and nodes map via the same hash function (e.g., Murmur3). ShardPay's ring state is versioned; routers cache ring v42 until the control plane publishes v43.
- Clockwise walk (successor) — First node at or after the key's hash owns the key. ShardPay's
ShardRouter.locate(accountId)performs binary search on sorted vnode positions — O(log vnodes) per lookup. - Virtual node (vnode) — Multiple ring positions per physical machine. ShardPay uses 150 vnodes per shard so a 16-shard cluster has 2,400 ring points — smooth load distribution even when physical nodes differ in CPU.
- Rebalance / migration — When vnodes move between physical shards, only key ranges between old and new successor boundaries copy. ShardPay's resharding job migrates ~1/16 of keys when growing 16→17 shards, not 15/16 as with
hash % 16. - Bounded loads — Extensions like Google's Jump Consistent Hash minimize migration further. ShardPay evaluated Jump Hash for even distribution; stayed with vnode ring for fine-grained hot-shard evacuation.
Modulo vs consistent hash
| Approach | Add 1 node to 16 | Hot-spot handling |
|---|---|---|
| hash(k) % 16 | ~15/16 keys remap | Poor — one key per bucket |
| Consistent hash + vnodes | ~1/17 keys remap per new node | Vnodes spread load |
hash(k) % 16
Add 1 node to 16~15/16 keys remapHot-spot handlingPoor — one key per bucketConsistent hash + vnodes
Add 1 node to 16~1/17 keys remap per new nodeHot-spot handlingVnodes spread load
Modulo hashing vs consistent hashing
Walkthrough: Growing from 16 to 17 shards
- Fleet runs 16 physical shards, 150 vnodes each (2,400 ring points).
- Capacity planning triggers add of
shard-17. - Control plane inserts 150 new vnodes for
shard-17at evenly spaced hashes. - Some key ranges that previously belonged to
shard-3now fall between new vnodes — those accounts migrate. - Migration service dual-writes to old and new shard for each migrating
accountIdduring cutover window. - Routers atomically switch to ring v43 after migration verification; ~6.25% of keys moved.
// ShardPay consistent hash ring (simplified)
public final class ConsistentHashRing {
private final SortedMap<Long, PhysicalShard> ring;
public PhysicalShard locate(String key) {
long hash = Hashing.murmur3_128().hashString(key, UTF_8).asLong();
SortedMap<Long, PhysicalShard> tail = ring.tailMap(hash);
Long vnode = tail.isEmpty() ? ring.firstKey() : tail.firstKey();
return ring.get(vnode);
}
}
Virtual nodes and hot-shard evacuation
Physical heterogeneity and traffic skew require more than naive one-node-one-point rings.
Advanced operations
- Vnode count tuning — More vnodes → smoother distribution, larger metadata. ShardPay uses 150 vnodes/shard; Cassandra often uses 256. Below 20 vnodes/shard, ShardPay saw 15% load skew in simulations.
- Weighted vnodes — High-memory shards get more vnodes. ShardPay's
shard-premiumhardware runs 300 vnodes vs 150 on standard — absorbs large tenants without directory overrides. - Handoff during migration — Client-side caching of shard location must invalidate on ring version change. ShardPay returns
X-Shard-Ring-Versionon every API response; clients refresh cache on mismatch. - Replication on ring — Successor walk can continue for replica placement (chord-style). ShardPay separates routing ring (account → shard) from intra-shard Raft replication — do not conflate the two layers.
- Hot key on ring — Consistent hashing spreads keys, not QPS. A single hot
accountIdmaps to one vnode regardless of vnode count — requires salting or splitting (see hot-spots article).
Example: Misconfigured single vnode per node
ShardPay's staging environment once deployed with 1 vnode per physical shard (16 points on the ring). Load tests showed shard-7 at 22% of traffic while shard-2 had 4% — hash clustering on a sparse ring. Increasing to 150 vnodes per shard brought skew under 2% without code changes — only metadata size grew.
System Design trackSee consistent hashing in system design context
Quick recall
Everything you need if you only revisit this box.
- Consistent hashing minimizes key remapping when nodes join or leave.
- Virtual nodes balance uneven segments — ShardPay uses 150 vnodes per physical shard.
- Prefer consistent hash over naive modulo for dynamic clusters; pair with migration tooling.
Test yourself
Answer these before moving on — recall is what makes it stick.