Skip to main content

Consistent Hashing

Pronunciation
kun-SIS-tunt HASH-ing
Updated 2 min read

Share this page

Send the link, quote the definition with a link back, or show it as a card on your own site.

https://softwaredictionary.org/terms/consistent-hashing

In short

Consistent hashing spreads keys across a changing set of servers so that adding or removing a server moves only a small share of the keys.

What is consistent hashing?

The simple way to pick a server for a key is hash(key) % number_of_servers. It spreads keys evenly, but when a server is added or removed, the result changes for nearly every key, so caches suddenly miss and data has to be moved everywhere. Consistent hashing, introduced in a 1997 paper for distributing web caches, avoids that.

Imagine the space of hash values as a ring. Each server is placed on the ring at the position of its own hash, and each key belongs to the first server found by moving clockwise from the key's position. When a server joins, it takes over only the keys between it and its neighbor; when one leaves, only its keys move to the next server. On average, about one key in n moves, where n is the number of servers.

With only a few servers the ring can be uneven, so each physical server is usually given many positions, called virtual nodes, which spreads load more evenly and lets stronger machines take more of it. Consistent hashing is used by DynamoDB and Cassandra to place data, by CDNs and caching layers to choose servers, and by load balancers that need the same client to keep reaching the same backend.

A common misconception is that consistent hashing balances load perfectly. It balances keys, not traffic: one very popular key, a hot key, can still overload the server that owns it. Systems add replication, splitting of hot keys or extra caching for those cases, and some use alternatives such as rendezvous hashing.

Key takeaways

  • Consistent hashing maps keys to servers on a hash ring.
  • Adding or removing a server moves only about 1/n of the keys.
  • Plain modulo hashing moves almost every key when servers change.
  • Virtual nodes give each server many positions for an even spread.
  • DynamoDB, Cassandra, CDNs and caches use it; hot keys still need care.

Example

A small hash ring with virtual nodes (Python)python
import bisect, hashlib

def h(value):
    return int(hashlib.md5(value.encode()).hexdigest(), 16)

class HashRing:
    def __init__(self, servers, vnodes=100):
        self.ring = sorted((h(f"{s}#{i}"), s) for s in servers for i in range(vnodes))
        self.keys = [position for position, _ in self.ring]

    def server_for(self, key):
        i = bisect.bisect(self.keys, h(key)) % len(self.ring)   # first server clockwise
        return self.ring[i][1]

ring = HashRing(["cache-a", "cache-b", "cache-c"])
print(ring.server_for("user:42"))
# Adding "cache-d" later moves only about a quarter of the keys.

Readers ask

Why not just use hash(key) modulo the number of servers?

Because changing the number of servers changes the result for almost every key, so nearly all cached or stored data ends up on the wrong server at once. Consistent hashing limits the move to a small fraction.

What are virtual nodes?

Extra positions on the hash ring for each physical server. Using many of them per server evens out the distribution of keys and lets more powerful servers take a bigger share.

Where is consistent hashing used?

In distributed databases such as Cassandra and DynamoDB to decide which nodes store which data, in distributed caches, in CDNs and in load balancers that keep a client on the same server.

See also

Spotted a mistake or something missing on this page?Suggest an edit

Read a random page
Open today's review
Switch to the dark theme
Read this page in Türkçe

More

Settings