skip to content

Why is consistent hashing used as a load-distribution strategy instead of plain modulo hashing, and what problem does it specifically solve when backend servers are added or removed?

level: seniorimportance: should knowfreq 55%

answer

  1. modulo hashing: N changes -> almost all keys remap
  2. ring: hash both servers and keys
  3. only ~1/N keys move on topology change
  4. virtual nodes even out load + spread failure
  5. used for sharded caches/data, not generic stateless routing

basics

~20 s

Plain modulo hashing (key % number-of-servers) reshuffles almost every key's assigned server whenever a server is added or removed, which is disastrous for caches. Consistent hashing arranges servers and keys on a ring so that adding or removing one server only reassigns a small fraction of keys, not nearly all of them.

solid answer

~60 s

With plain modulo hashing, a key's assigned server is hash(key) mod N, where N is the server count; changing N (adding or removing even one server) changes the modulo result for almost every key, meaning almost every client's requests suddenly map to a different backend - catastrophic for anything relying on server-local state, like a cache, since it causes a near-total cache invalidation cascade back to the origin. Consistent hashing solves this by mapping both servers and keys onto positions on a fixed hash ring (0 to 2^32-1, say) and assigning each key to the next server clockwise from its position; adding or removing a server only affects the keys between it and its immediate neighbor on the ring, leaving the rest of the mapping untouched. Real implementations add virtual nodes - each physical server is hashed to many points on the ring - so that load stays evenly distributed even with few servers, since a single physical node's failure spreads its lost range across many other nodes rather than dumping it all on one neighbor.

go deeper

for a junior

Should have a rough sense that adding a server shouldn't scramble everyone's assignment, without needing ring mechanics.

for a middle

Should explain the modulo-hashing reshuffle problem and that consistent hashing bounds the number of keys that move.

for a senior

Should describe the ring mechanism, the ~1/N reshuffle bound, and why virtual nodes are needed in practice.

for a principal

Should place consistent hashing correctly relative to ordinary load-balancing algorithms (stateful key ownership vs stateless request distribution) and reason about its use in real sharded systems (Dynamo-style stores, cache clients) and its limits (hot keys, added complexity when statelessness would suffice).

## Modulo hashing and its one severe weakness Modulo hashing is the naive approach to mapping a key (a user ID, a cache key, a client identifier) onto one of N backend servers: compute `hash(key) mod N` and route to that server index. It's simple and gives a uniform distribution for a fixed N, but it has a severe weakness the moment N changes. Because the modulo operation depends on the exact value of N, adding or removing even a single server changes the divisor, which changes `hash(key) mod N` for the overwhelming majority of keys - not just the keys that logically belong to the added or removed server, but essentially all of them, since the entire modular arithmetic shifts. - For a **stateless** load balancer picking any interchangeable backend, this churn is mostly harmless, since any backend can serve any request. - But the moment servers hold **key-dependent state** - most commonly, a cache where server X holds the cached value for keys that hash to X - this remapping is catastrophic: on a scaling event, nearly every key now maps to a different, cold server, so nearly every subsequent request is a cache miss, and all of that traffic floods back to the origin database or service simultaneously. This is sometimes called a 'cache stampede' or 'cold cache cascade,' and it can be severe enough to take down the origin that the cache was there to protect, at exactly the moment (a scaling event, often triggered by high load) when the system is least able to absorb it. ## The ring, and the locality of change it buys Consistent hashing was designed specifically to bound this blast radius. The core idea is to hash both the servers and the keys into the same abstract space - conventionally visualized as a ring of hash values from 0 to some maximum (e.g. 2^32 - 1) - by hashing each server's identifier (its hostname or address) to get its position on the ring, and hashing each key the same way to get its position. A key is then assigned to the first server encountered walking clockwise from the key's position on the ring. The crucial property this produces is **locality of change**: - when a new server is added, its position on the ring only intercepts the range of keys between it and the previous server going counter-clockwise - only those keys move to the new server, and every other key's assignment is completely unaffected; - symmetrically, when a server is removed, only the keys it owned move to its next clockwise neighbor; every other server's assignments are untouched. Instead of nearly 100% of keys reshuffling on a topology change, only roughly 1/N of the keys move, where N is the number of servers - a dramatically smaller, and more importantly, bounded and predictable disruption. ## Why real implementations add virtual nodes A naive version of consistent hashing has its own problem though: with only a handful of physical servers hashed to single points on the ring, the ranges between them can be quite uneven purely by chance, giving some servers a much larger share of the key space (and therefore load) than others, and when a server is removed, its entire range dumps onto exactly one neighbor rather than spreading out. Real implementations fix this with **virtual nodes** (sometimes called 'vnodes'): instead of hashing each physical server to one point on the ring, each physical server is hashed to many points (commonly 100-200 per server), scattered across the ring under different virtual identifiers, and a key is still assigned to whichever virtual node it reaches first, which then maps back to the owning physical server. This does two things: 1. it evens out the load distribution because the law of large numbers smooths out the ranges when there are many points per server rather than one; 2. it means that when a physical server fails, its many small ranges are spread across many different neighboring physical servers rather than dumping all of that load onto a single one. ## Where it is actually used The canonical real-world use is exactly the cache-sharding scenario described above. These all use consistent hashing to decide which node owns a given key, precisely because minimizing reshuffling on node membership changes is the whole point: - systems like Memcached client libraries (e.g. `libketama`), - Amazon's Dynamo (and derivatives like Cassandra and Riak for partitioning data across nodes), - and CDN request routing. ## How it differs from an ordinary balancing algorithm It's worth being precise about where this fits relative to ordinary load-balancer algorithms like round-robin or least-connections. | Mechanism | The question it answers | |---|---| | round-robin, least-connections | which backend handles the next request when any backend can serve it equally well (stateless request distribution) | | consistent hashing | which backend owns a given piece of state or data (stateful key-to-node assignment) | Consistent hashing is specifically valuable at a load balancer or client layer when you need requests for the same key to consistently land on the same backend - not merely for session affinity/locality (where losing the mapping just costs a cache miss) but for correctness in sharded systems (where the 'right' node genuinely owns that data). A load balancer offering 'consistent hashing' as a balancing mode is really offering a way to get affinity-like key-to-backend stability that degrades gracefully, rather than catastrophically, as the backend pool changes size.

  • Why do virtual nodes specifically help when a single physical server fails, beyond just improving average load balance?
    Without virtual nodes, a failed server's entire key range falls onto just one neighboring server on the ring, potentially doubling that single neighbor's load overnight. With virtual nodes, the failed server's range is actually many small scattered ranges, each falling onto a different neighboring physical server, so the lost capacity is absorbed as a small proportional increase spread across many servers instead of a large spike on one.
  • In what scenario would consistent hashing not actually help, even though it sounds like the right tool?
    If backends are genuinely stateless and interchangeable - any server can serve any request with identical cost and no cache/locality benefit - consistent hashing adds complexity (ring maintenance, virtual node bookkeeping) without a corresponding benefit, since a simpler algorithm like least-connections already balances load well and there's no state whose remapping cost needs bounding.
  • How does the choice between consistent hashing and simple sharding-by-range (e.g. by key prefix) compare for handling growth?
    Range-based sharding needs an explicit, often manually managed split when a range gets too hot, and rebalancing can be an all-or-nothing operation on that range's data; consistent hashing automatically and incrementally redistributes only the affected slice of key space when a node is added, without needing a coordinator to decide where to split, though it can still result in uneven 'hot key' ranges that neither approach solves on its own.

Modulo hashing is like renumbering every locker in a school whenever one locker bank is added or removed - everyone's locker number changes and nobody can find their stuff. Consistent hashing is like adding lockers only at the end of a long hallway - only the people nearest the new lockers get reassigned, and everyone else keeps their same locker.

saying these in an interview costs you the question

  • Thinks modulo hashing and consistent hashing perform equally well on scaling events
  • Can't explain why adding one server under modulo hashing reshuffles almost all keys
  • Doesn't mention virtual nodes when discussing consistent hashing's practical implementation
  • Confuses consistent hashing with session affinity/sticky sessions
  • Believes consistent hashing is only relevant to caches and never to sharded databases

context