Distributed Systems

Consistent Hashing Explained

How consistent hashing works, why it beats modulo hashing when servers change, how virtual nodes balance load, and where it’s used in caches and databases.

Abstract ring of glowing nodes representing a consistent hashing ring
Illustration: Backend Architect / AI-generated.

Key takeaways

  • With modulo hashing, adding one server remaps almost every key; consistent hashing moves only a small fraction.
  • Keys and servers are placed on a ring; each key belongs to the next server clockwise.
  • Virtual nodes smooth out uneven load and make rebalancing gentler.
On this page

When you spread data across many cache or database nodes, you need a rule for which node owns which key. Consistent hashing is the rule most large systems use, because it keeps working smoothly when nodes are added or removed.

The problem with modulo hashing

The simple approach is node = hash(key) % N. It distributes keys evenly, until N changes. Add a fifth server to four and almost every key maps to a different node. For a cache, that means a sudden flood of misses; for a database, it means moving nearly all your data.

How consistent hashing works

  1. Imagine a circle (the ring) of hash values.
  2. Hash each node to a position on the ring.
  3. Hash each key to a position on the ring.
  4. A key belongs to the first node clockwise from its position.

When a node joins, it takes over only the keys between it and the previous node. When a node leaves, only its keys move to the next node. On average, only about 1/N of keys move.

Virtual nodes

With a handful of physical nodes, positions on the ring can be uneven, leaving some nodes with far more keys. The fix is virtual nodes: each physical node is placed on the ring many times under different hashes.

  • Load spreads more evenly.
  • When a node fails, its keys scatter across many nodes rather than overloading one neighbour.
  • More powerful machines can be given more virtual nodes.
import bisect, hashlib

class Ring:
    def __init__(self, nodes, vnodes=100):
        self.ring = sorted((self.h(f"{n}#{i}"), n) for n in nodes for i in range(vnodes))
        self.keys = [k for k, _ in self.ring]

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

    def node_for(self, key):
        i = bisect.bisect(self.keys, self.h(key)) % len(self.ring)
        return self.ring[i][1]

Where it’s used

Replication on the ring

To keep copies of data, a key is stored on its owner and the next few distinct nodes clockwise. If one fails, the others still serve it, a pattern tied closely to the trade-offs in the CAP theorem.

Frequently asked questions

How many virtual nodes should each server have?

Often somewhere between tens and a few hundred. More gives smoother balance at the cost of a larger ring to search.

Does consistent hashing guarantee perfect balance?

No, but with virtual nodes it gets close, and it minimises data movement when nodes change.

What hash function should I use?

Any fast function with good distribution. Cryptographic strength isn’t required.

Sources

  1. Karger et al. — Consistent Hashing and Random Trees (1997)
  2. Martin Kleppmann — Designing Data-Intensive Applications

Every article is edited by a human and checked against our editorial policy. Spotted a mistake? Tell us.

Keep reading

Distributed Systems

How CDNs Work

How CDNs work: edge locations, request routing, cache keys, Cache-Control headers, invalidation, origin shielding, security features and edge compute.

4 min read