High Level Design
Consistent Hashing
How consistent hashing solves the massive-rehashing problem in distributed systems — hash rings, virtual nodes, load imbalance, and real-world use cases.
In distributed systems, data must be distributed across multiple servers. When the cluster size changes, the naive approach causes almost all keys to remap — consistent hashing solves this.
Naive Approach#
- Use a modulo hash: serverIndex = hash(key) % N to assign keys to servers.
- Works when pool size N is fixed and data distribution is even.
Problem: When N changes (node added or removed), almost all keys remap to different nodes — massive data movement, cache misses, and high rehashing cost.
We need a hashing mechanism that minimizes re-distribution of keys when nodes change.
Consistent Hashing#
A hashing technique that ensures:
- Only a small fraction of keys remap when nodes join/leave.
- Keys are distributed more evenly across servers.
When a hash table is resized using consistent hashing, only K/n keys need to be remapped on average, where K = number of keys and n = number of slots.
The Technique — Step by Step#
| Step | Depiction |
|---|---|
| Hash Ring: Represent the hash space as a ring (0 to 2³² – 1). | ![]() |
| Hash servers: Map servers to the ring using their IP or name via a hash function. | ![]() |
| Hash Keys: Both keys and nodes are placed on the ring using the same hash function. | ![]() |
| Each key is assigned to the next node clockwise on the ring. | ![]() |
Handling Node Changes#
Adding a Node#
- The new server is hashed to find its position on the ring.
- It takes over responsibility for keys between itself and its predecessor.
- Only keys that map to the new server need to move.
In the figure below, after server 4 is added, only key0 needs to remap.

Removing a Node#
- The removed server's keys are reassigned to the next server clockwise.
- Only keys assigned to the removed server need to move.
In the figure below, after server 1 is removed, only key1 needs to remap.

Issues with Basic Consistent Hashing#
1. Hot Spot Key Problem#
When a server is removed, the next server takes over all its keys — load spikes.

2. Load Imbalance#
If servers land unevenly on the ring, some get far more keys than others.

Solution: Virtual Nodes#
Each physical server is assigned multiple positions on the ring → better key distribution.
- When a server is added/removed, only a fraction of its virtual nodes are affected.
- Key assignment works the same: move clockwise to the next (virtual) node.

Benefits#
| More virtual nodes → | More balanced key distribution |
|---|---|
| Standard deviation of load | Gets smaller as virtual nodes increase |
| Weighting | Assign more virtual nodes to higher-capacity servers |
| Trade-off | More storage needed to track virtual node metadata |
Finding Affected Keys to Remap#
When a server is added:
- Start from the new server's position, go anticlockwise until another server is found.
- All keys in that range remap to the new server.
When a server is removed:
- Start from the removed server's position, go anticlockwise until another server is found.
- All keys in that range remap to the next server clockwise.
Design Flow#
- Client sends a key-based request.
- Key is hashed → position found on hash ring.
- Lookup algorithm finds the nearest node clockwise.
- Request is routed to that server.
Production Considerations#
- Node membership: Use gossip protocols, Zookeeper, or etcd to keep membership updated.
- Replication: Store each key on the next N clockwise nodes for fault tolerance.
- Scaling: Adding/removing nodes causes minimal reshuffling — only nearby keys move.
Example Use Cases#
- Distributed caches (Memcached, Redis Cluster)
- Distributed databases (Cassandra, DynamoDB, Riak)
- CDNs — cache routing to the nearest edge node
- Sharded key-value stores
- Load balancers — consistent request-to-server mapping
Real-World Examples#
- Amazon DynamoDB — uses consistent hashing for data distribution and replication. Reference
- Apache Cassandra — uses consistent hashing to distribute data across cluster nodes. Reference
- Maglev (Google) — Google's software load balancer uses consistent hashing for request distribution. Reference
Gist#
Consistent hashing solves massive rehashing by mapping both nodes and keys onto a hash ring. When nodes join or leave, only K/n keys remap. Virtual nodes ensure even load distribution. Widely used in distributed caches, databases, and CDNs.



