Consistent hashing is a distributed partitioning algorithm that maps both data keys and physical storage nodes to positions on a circular hash ring. When a storage node is added or removed, consistent hashing requires reallocating only a small fraction of keys proportional to 1/N, where N is the total number of nodes. Distributed data platforms like Apache Cassandra, Amazon DynamoDB, and Memcached clusters use consistent hashing to avoid the catastrophic data reshuffling caused by naive hash modulo strategies.
Ring topology and key remapping
Under a naive modulo hash scheme, a record key is routed using the formula hash(key) mod N. If the cluster size changes from 10 nodes to 11 nodes, almost every existing key maps to a completely different node, invalidating distributed caches and forcing full cluster data migrations:
- The circular ring: Consistent hashing uses a uniform 32-bit or 64-bit integer range that wraps around from 0 to 2^64 - 1. Physical nodes are hashed by their IP address or hostname to claim fixed positions on the circle.
- Key assignment: When a write arrives, the system hashes the record partition key to locate its point on the ring. The key travels clockwise along the circumference until it encounters the first node, which becomes its coordinator or primary owner.
- Fault tolerance: If a node fails, its keys simply shift clockwise to the next adjacent node. Peer nodes remain untouched, meaning only keys previously assigned to the dead node require migration.
[Node A: pos 0]
. .
[Node C: pos 66] [Node B: pos 33]
. .
[Ring 0 -> 2^64-1]
Key hashes to pos 15 -> routes clockwise to Node B.
Add Node D at pos 20 -> only keys between 0 and 20 move to Node D!Physical hardware limitations require virtual tokens to avoid uneven distribution:
Virtual nodes and hotspot mitigation
In a basic ring setup, randomly placing three or four nodes leads to uneven ring segments, causing one node to receive half of the cluster traffic.
Distributed engines solve this using virtual nodes (vnodes). Each physical server is assigned hundreds of distinct token positions scattered across the ring. If physical server A has 256 virtual nodes, its responsibilities are dispersed evenly across the entire keyspace. When server A is decommissioned, its 256 virtual segments are absorbed by dozens of different physical peers rather than dumping the entire burden onto a single adjacent neighbor.