PrepZone Logo
PrepZone

Consistent Hashing with Virtual Nodes

Add or remove nodes without remapping the entire key space — ShardPay's shard ring.

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.
Hash ring
Node Avn1, vn2
Node Bvn3, vn4
Node Cvn5, vn6

Keys hash onto the ring; walk clockwise to find owner

Virtual nodes spread keys evenly. Adding a node moves only adjacent key ranges.

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

ApproachAdd 1 node to 16Hot-spot handling
hash(k) % 16~15/16 keys remapPoor — one key per bucket
Consistent hash + vnodes~1/17 keys remap per new nodeVnodes spread load
  • hash(k) % 16

    Add 1 node to 16~15/16 keys remap
    Hot-spot handlingPoor — one key per bucket
  • Consistent hash + vnodes

    Add 1 node to 16~1/17 keys remap per new node
    Hot-spot handlingVnodes spread load

Modulo hashing vs consistent hashing

Walkthrough: Growing from 16 to 17 shards

  1. Fleet runs 16 physical shards, 150 vnodes each (2,400 ring points).
  2. Capacity planning triggers add of shard-17.
  3. Control plane inserts 150 new vnodes for shard-17 at evenly spaced hashes.
  4. Some key ranges that previously belonged to shard-3 now fall between new vnodes — those accounts migrate.
  5. Migration service dual-writes to old and new shard for each migrating accountId during cutover window.
  6. Routers atomically switch to ring v43 after migration verification; ~6.25% of keys moved.
Java
// 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-premium hardware 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-Version on 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 accountId maps 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.

Quick recall

Everything you need if you only revisit this box.

  1. Consistent hashing minimizes key remapping when nodes join or leave.
  2. Virtual nodes balance uneven segments — ShardPay uses 150 vnodes per physical shard.
  3. 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.