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.

August 10, 2026

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#

StepDepiction
Hash Ring: Represent the hash space as a ring (0 to 2³² – 1).hash ring
Hash servers: Map servers to the ring using their IP or name via a hash function.hash servers
Hash Keys: Both keys and nodes are placed on the ring using the same hash function.hash keys
Each key is assigned to the next node clockwise on the ring.finding servers

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.

adding server

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.

removing server


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.

hotspot problem

2. Load Imbalance#

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

load imbalance


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.

virtual nodes

Benefits#

More virtual nodes →More balanced key distribution
Standard deviation of loadGets smaller as virtual nodes increase
WeightingAssign more virtual nodes to higher-capacity servers
Trade-offMore storage needed to track virtual node metadata

Finding Affected Keys to Remap#

When a server is added:

  1. Start from the new server's position, go anticlockwise until another server is found.
  2. All keys in that range remap to the new server.

When a server is removed:

  1. Start from the removed server's position, go anticlockwise until another server is found.
  2. All keys in that range remap to the next server clockwise.

Design Flow#

  1. Client sends a key-based request.
  2. Key is hashed → position found on hash ring.
  3. Lookup algorithm finds the nearest node clockwise.
  4. 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.