πŸ—“οΈ 29042026 1515
πŸ“Ž #distributed_systems #sharding

CONSISTENT HASHING

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​

ABSTRACT

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​

SystemNotes
DynamoDBOriginal Dynamo paper popularised it. Vnodes called "tokens".
CassandraSame 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.
Riak64-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, not 1/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 N is the wrong default, since adding one node remaps (N-1)/N of keys, not 1/N. The naive scheme is the alternative this exists to replace.

References​

  • Karger et al., "Consistent Hashing and Random Trees" (1997)
  • DeCandia et al., "Dynamo: Amazon's Highly Available Key-value Store" (SOSP 2007)