ποΈ 29042026 1515
π #distributed_systems #sharding
Sharding scheme that minimises data movement when nodes join or leave the cluster. The basis for distributed caches and Dynamo-style stores β the alternative ("just
hash mod N") falls apart on the first capacity change.
The Problem It Solvesβ
Naive sharding: node = hash(key) mod N. Works fine β until N changes.
Add or remove one node, almost every key rehashes to a new node. For a cache, that's a stampede onto the origin. For a stateful store, that's TB-scale data movement.
You want: when N changes by 1, only ~1/N of keys move.
The Mechanismβ
Map both keys and nodes to the same hash space (typically [0, 2^32)), arranged as a ring. Each key is owned by the first node clockwise from its hash position.
hash space (ring, 0 β 2^32 β 0)
βββββββββββββββββββ
β βnodeA β
β βkeyX β
β βnodeB β
β βkeyY β
β βnodeC
βββββββββββββββββββ
keyX β walk clockwise β nodeB
keyY β walk clockwise β nodeA
When nodeB leaves: only keys in the arc (nodeA, nodeB] reassign β they walk further clockwise to nodeC. Keys owned by nodeA and nodeC are untouched.
Add a node anywhere on the ring: it steals an arc from one neighbour. Same locality property.
Virtual Nodes (vnodes)β
Naive ring placement gives uneven arcs β uneven load. With only 3 nodes, one node may own 60% of the ring by chance.
Fix: each physical node registers many virtual positions on the ring (e.g. 100β200 vnodes per physical node). Hash node_id || replica_index for each. With enough vnodes, arc sizes converge to roughly equal.
Bonus: when a node fails, its load distributes across all remaining nodes (because its vnodes are scattered), not piled onto its single clockwise neighbour.
Where It's Usedβ
| System | Notes |
|---|---|
| DynamoDB | Original Dynamo paper popularised it. Vnodes called "tokens". |
| Cassandra | Same lineage. Default num_tokens=16 per node (was 256 in older versions). |
| Memcached (clients) | Client-side consistent hashing (e.g. ketama) routes keys to servers without coordination. |
| Riak | 64-bit ring divided into fixed partitions assigned to nodes. |
Redis Cluster does NOT use consistent hashing β it uses 16384 fixed slots assigned to nodes (see redis_cluster). Different scheme, similar goal: bounded data movement on resharding.
Common Pitfallsβ
- Data movement when adding a node is
~1/(N+1)of keys, not1/N. The exact fraction matters for capacity planning. - Replication on the ring β after finding the owner, replicas live on the next R-1 nodes clockwise. This is how Dynamo-style systems get N=3 replication "for free".
- Hot keys still hot β consistent hashing balances key distribution, not access frequency. A celebrity's user_id still hammers one node. Mitigation: separate hot-key handling layer or per-key replication.
- Hash function quality matters β MD5/MurmurHash3 are fine; a weak hash gives clumpy rings even with vnodes.
hash mod Nis the wrong default, since adding one node remaps(N-1)/Nof keys, not1/N. The naive scheme is the alternative this exists to replace.
Relatedβ
- redis_cluster β fixed-slot variant of the same problem.
- quorum_reads_writes β replication on the ring with N/W/R.
- cap_theorem β Dynamo-style systems built on consistent hashing are typically AP.
Referencesβ
- Karger et al., "Consistent Hashing and Random Trees" (1997)
- DeCandia et al., "Dynamo: Amazon's Highly Available Key-value Store" (SOSP 2007)