Match a job Paths Subjects Questions Quizzes Pricing

Consistent Hashing

Distributing data across nodes with minimal reshuffling

Overview Read

Consistent Hashing

When you have data spread across multiple servers, you need a way to answer: which server is responsible for a given key? The naive approach — server = hash(key) % number_of_servers — has a fatal flaw: adding or removing a server changes the modulus, causing almost every key to remap to a different server. In a distributed cache, this means a near-total cache miss storm. In a distributed database, it triggers massive data migration.

Consistent hashing solves this by arranging servers and keys on a circular "ring" so that adding or removing a node only remaps the keys that were previously owned by that node — typically 1/N of total keys, where N is the number of servers.


The Hash Ring

Imagine a circle with values from 0 to 2³²-1 (or 0 to 2⁶⁴-1 for a 64-bit hash space). Both servers and keys are hashed to positions on this ring using the same hash function.

                    0 / 2^32
                   /
              S3 ●             ● S1
                               
                         K2  
                    K1 ●         ● K3
                               
              S2 ●             
                               
                   2^32 / 2

Routing rule: each key is assigned to the first server encountered when moving clockwise around the ring from the key's position.

  • K1 → S1 (first server clockwise from K1's position)
  • K2 → S1
  • K3 → S3 (wraps around the top)

Adding a Server

When S4 is added at a position between K1 and S1, only K1 remaps to S4. All other keys remain with their current servers. Only ~1/N keys are affected.

Removing a Server

When S2 is removed, the keys that were assigned to S2 now map to the next server clockwise (S3). All other keys are unaffected.


Virtual Nodes

In the basic ring, if servers have different capacities (e.g., one server has 4× the RAM), or if the hash function places servers unevenly, some servers handle far more keys than others (hot spots).

Virtual nodes (vnodes) solve this: each physical server is assigned multiple positions on the ring, distributed evenly. A server with more capacity gets more virtual nodes.

Physical servers: S1 (high-end), S2 (standard), S3 (standard)
Virtual nodes:    S1 → positions [12, 45, 78, 112, 145, 178] (6 vnodes)
                  S2 → positions [20, 55, 90, 125]           (4 vnodes)
                  S3 → positions [30, 65, 100, 135]          (4 vnodes)

With enough virtual nodes (Cassandra uses 128–256 per node), load is distributed much more evenly, and adding/removing a physical server redistributes load from/to all other servers rather than just its immediate neighbors.


Pro content

Sign up free, then start a 14-day Pro trial — no card needed.

We use cookies for product analytics to improve OmniAtlas. See our Privacy Policy.