Consistent Hashing: Sharding Without the Reshuffle

Add a cache node with plain modulo hashing and almost every key moves. Consistent hashing moves only one arc. How the hash ring and virtual nodes work.

You have eight cache nodes, and you spread keys across them with hash(key) % 8. It works: reads are fast, load is even, everyone is happy. Then traffic grows and you add a ninth node. The moment you change that 8 to a 9, almost every key you have ever cached maps to a different node. Every lookup misses, and all nine nodes ask the database for data they used to hold. The database that sat comfortably behind a cache is suddenly taking the full read load with no warmup. Adding capacity is what knocked the system over.

That failure has a name and a fix. The fix is consistent hashing, and it is the reason distributed caches, sharded databases, and load balancers can change their member list without a stampede.

This post is for backend and infrastructure engineers who route keys across a pool of nodes: a Redis or Memcached tier, shards of a store, or a set of worker processes. You want adding or removing one node to move a small, predictable slice of the keys instead of all of them. I will walk through why modulo breaks, what the ring actually is, a working implementation in about thirty lines, and the ways it still bites in production.

Why modulo hashing falls apart

The problem with hash(key) % N is that N is baked into the answer. Change N and you change the answer for nearly every key, because the old and new results almost never agree.

Take the jump from eight nodes to nine. A key stays on the same node only when hash(key) % 8 equals hash(key) % 9. How often does that happen? Both remainders are fully decided by hash(key) % 72, because 72 is the least common multiple of 8 and 9. Of those 72 possible remainders, only 8 give matching results. That is one in nine. So roughly eight of every nine keys land on a different node the instant you scale out.

The math is the same in the other direction. Losing a node to a crash reshuffles just as violently, at the worst possible moment.

Modulo also assumes your nodes are numbered 0..N-1 with no gaps. Real fleets are not like that. A node dies, you replace it with one that has a different address, and now you are renumbering the world. What you want is a scheme where each node claims a fixed region of the key space, independent of how many other nodes exist. That is exactly what consistent hashing gives you.

The hash ring

Consistent hashing was introduced by Karger et al. in 1997 for exactly this “relieving hot spots” problem. It starts with one idea: stop treating the hash output as a bucket index and treat it as a position. Hash everything, keys and nodes alike, into the same large numeric space, say 64-bit integers. Then bend that space into a circle, so the largest value wraps back around to zero.

Each node hashes to some point on the ring, and so does each key. A key belongs to the first node you reach going clockwise from the key’s position. That rule is the whole system.

Two ring diagrams side by side. On the left, three nodes A, B and C sit on a circle and four keys k1 to k4 each connect by a dashed arrow to the first node clockwise from them; k2 lands on B. On the right the same ring has a fourth node D inserted on the arc between A and B; only k2 changes owner, moving from B to D, while k1, k3 and k4 stay exactly where they were. The caption reads: adding a node moves one arc, not the whole map.

The payoff is on the right side of that picture. When you add node D, it lands on the ring at one spot and takes over the arc immediately counter-clockwise of it. Those are the keys that used to walk past that stretch to reach B. Every other key is untouched, because its clockwise walk has not changed. On average, adding a node to an N-node ring moves 1/(N+1) of the keys, and it only disturbs the two nodes adjacent to the newcomer. Compare that to the eight-of-nine carnage from modulo.

Here is the ring as real code. Lookups stay fast because the node positions live in a sorted array, and a binary search finds the first position clockwise of a key:

import bisect
import hashlib

class HashRing:
    def __init__(self, nodes=None, vnodes=100):
        self.vnodes = vnodes
        self._positions = []   # sorted ring positions
        self._owner = {}       # position -> node name
        for node in nodes or []:
            self.add(node)

    def _hash(self, key: str) -> int:
        digest = hashlib.md5(key.encode()).digest()
        return int.from_bytes(digest[:8], "big")   # 64-bit position

    def add(self, node: str) -> None:
        for i in range(self.vnodes):
            pos = self._hash(f"{node}#{i}")
            self._owner[pos] = node
            bisect.insort(self._positions, pos)

    def remove(self, node: str) -> None:
        for i in range(self.vnodes):
            pos = self._hash(f"{node}#{i}")
            del self._owner[pos]
            self._positions.remove(pos)

    def get(self, key: str) -> str:
        if not self._positions:
            raise KeyError("ring is empty")
        pos = self._hash(key)
        idx = bisect.bisect(self._positions, pos) % len(self._positions)
        return self._owner[self._positions[idx]]

The interesting line is in get. bisect.bisect returns the index of the first position greater than the key’s hash. The % len wraps you back to the start of the array when the key sits past the last node, which is the “wrap around the circle” step. Lookups are O(log n) in the number of positions, and n here is small: it is the node count, not the key count. The ring never stores the keys themselves. It only answers “which node owns this hash,” and the node holds the data.

Balance: where the naive ring goes wrong

You may have spotted the vnodes parameter in the code, which I have not explained yet. Without it, the ring has a real problem, and it is worth seeing why before trusting it in production.

If each node is a single point on the ring, the arcs between points are whatever the hash function happened to produce. Hash functions do not deal out even gaps. With three nodes, one can easily own half the ring while another owns a sliver. That imbalance is load: the node with the big arc gets a proportional share of the traffic and the memory pressure. You did not choose which node draws the short straw. The hash did.

Two ring diagrams comparing key-space ownership. On the left, three nodes A, B and C each sit at one point on the ring, and the colored arcs they own are visibly unequal, with A owning far more of the ring than B. On the right, the same three nodes are each placed at three points around the ring for nine points total; the colored arcs are now interleaved and nearly equal in size. The caption reads: virtual nodes even out the arcs.

The fix is virtual nodes. Instead of placing each physical node on the ring once, you place it many times (a hundred is a common default) by hashing node#0, node#1, and so on. Each of those points owns a small arc, and one physical node’s arcs are scattered all around the ring. With enough points, the law of large numbers takes over: the many small random arcs average out, and every node ends up owning close to the same total fraction. That is what the for i in range(self.vnodes) loops in add and remove do.

Virtual nodes buy a second thing, which matters more than balance. When a node fails, its arcs are spread across the ring, so its load spills onto many different neighbors. Without virtual nodes, the whole load lands on the one node that happened to sit clockwise of the failed one. That doubles the load on exactly one survivor, often enough to knock it over too, and now you have a cascade. Spreading the departed node’s share across the whole fleet is the difference between a shrug and an outage.

Where it bites

The ring is elegant on a whiteboard. These are the things I have learned to check before relying on it.

Even ownership of the key space is not even ownership of the traffic. Consistent hashing balances keys, on the assumption that keys cost roughly the same. One hot key (a celebrity account, a viral document, the config blob every request reads) still lands on exactly one node and hammers it. No number of virtual nodes helps, because that key is a single point on the ring. Hot keys need their own answer: replicate them, cache them client-side, or special-case them. The ring solves distribution, not skew.

The hash function has to spread bits, not avalanche cryptographically. You want uniform output, and you want it cheap, because you call the hash on every request. I used MD5 above because it is everywhere and its distribution is fine for this. Its brokenness as a cryptographic hash is irrelevant when you are only bucketing, not defending against an adversary. In a hot path I would reach for a non-cryptographic hash like xxHash or MurmurHash instead: same uniformity, a fraction of the cost. Do not use Python’s built-in hash() for a ring that must agree across processes. It is salted per interpreter run by default, so two machines will disagree about where a key lives.

Even with virtual nodes, load is only balanced in expectation. A hundred points per node keeps most nodes within a tolerable band, but the tail can still run hot, and randomness guarantees some node draws more than its share. If you need a hard ceiling rather than a statistical one, Google’s consistent hashing with bounded loads (paper, 2016) adds a capacity cap. A key that would land on a full node keeps walking clockwise to the next node with room. You give up a little of the “keys never move unnecessarily” purity in exchange for a firm bound on how overloaded any one node can get.

Rebalancing still costs something. “Only 1/N of keys move” is a real win, but it is not zero. Those keys still miss on their new node and get refetched from the source. When you add or drain a node, that fraction of your cache goes cold at once. Scale during a quiet window if you can, and make sure the backing store can absorb the miss spike. This is the same instinct behind warming a cache before a deploy rather than after.

When the ring is more than you need

Consistent hashing earns its complexity when nodes join and leave often and you cannot renumber them freely. There is a lighter option if your buckets are simply numbered 0..N-1 and you only ever grow the count (think shards of a data file, not a fleet of churning cache boxes). Jump consistent hash (Lamping and Veach, 2014) does the same job, moving the minimum number of keys when N changes, in five lines of arithmetic. It needs no ring, keeps no stored state, and balances better.

Its catch is right there in the constraint: buckets must be a contiguous range of integers, so you cannot remove node number 3 and keep 4 through 8. Pick the tool by how your nodes are named and how they come and go, not by which algorithm is more famous.

Where this shows up in my work

I reach for consistent hashing whenever a request has to find the same node twice. In Archi, the RAG (retrieval-augmented generation) copilot for CMS operations at CERN, the obvious place is the cache tier in front of retrieval and generation. It is the same idea I wrote about in semantic caching for LLM apps. You want a given query to keep hitting the node that already holds its embedding or its cached answer. You also want to add a cache box under load without cold-starting the whole tier.

It is also the mechanism hiding inside systems I lean on daily. OpenSearch routes a document to a shard by hashing its routing value. That is why changing the primary shard count forces a full reindex rather than a quiet resize: the partitioning function moved under it.

None of this is new. The 1997 paper that named it was aimed at web caches. It has lasted because of the property it guarantees: change the membership, and only a small, bounded set of keys moves. You need that every time a distributed system has to grow, shrink, or survive a node dying at 3am. Understanding the ring is the difference between adding a node as a routine operation and adding a node as an incident.


Diagrams by M. Hassan Ahmed, released under CC0. No external image was used for this post; the figures are original work by the author.