Intermediate System design concept · Distributed Systems & Data · 30 mins read

Consistent Hashing

Place keys on servers so that adding or removing a server moves only a small fraction of the data.

The Hash Ring

Why modulo placement breaks when servers change, and how a hash ring limits movement to about 1/N of the keys.

Intuition

With four cache servers and placement hash(key) % 4, adding a fifth server changes the target for about 80% of keys. A cache suddenly misses on most requests, and the database behind it gets a traffic spike exactly when you were trying to add capacity.

Consistent hashing, introduced for web caching in 1997, is the standard way to place data in clusters whose membership changes. It underpins Dynamo-style databases, cache clients and several load balancers.

Mental Model

Hash every server's ID to a position on a circle of, say, 2^32 points. Hash each key onto the same circle. A key belongs to the first server at or after its position, going clockwise. Adding a server only takes keys from the arc between it and the previous server; removing a server hands its arc to the next one. On average, one of N servers changing moves about 1/N of the keys. Think of it like: Clock positions: each hour mark is a server, and every minute belongs to the next hour mark. Add a mark at half past and only the minutes between it and the previous mark change owner.

Building Blocks

  • Hash space as a circle: A large output range (for example 0 to 2^32 − 1) wrapped so the highest value is next to zero.
  • Server positions: Each server is placed at hash(server_id); positions are stable as long as the ID is stable.
  • Clockwise successor lookup: A key's owner is the first server at or after hash(key). With servers kept in a sorted array, this is a binary search: O(log N).
  • Local membership view: Every client or router holds the same list of servers, so they all compute the same owner without a central lookup.

Definitions

Consistent hashing
A placement scheme where a membership change moves only about K/N of K keys.
  • Modulo hashing moves most keys on any change.
Successor
The first server clockwise from a key's position, which owns that key.
Remapping
Keys changing owner after a membership change.
  • In a cache this means misses.
  • In a database it means data must be moved.

Patterns

  • Client-side ring — When clients talk directly to cache or storage servers.
  • Ring in the router — When requests for the same key should reach the same backend.

Strategies

  • Hash stable IDs, not addresses When: When servers can be replaced or rescheduled. How: Hash a logical node name rather than an IP, so a restart on a new IP keeps the same position and the same keys. Example: Hash 'cache-07' instead of '10.0.3.41'.
  • Change membership gradually When: When adding capacity to a busy cache. How: Add one server at a time and let its share warm up before adding the next, so only one arc is cold at once. Example: Scale a cache pool from 8 to 12 servers over an hour instead of all at once.

Alternatives to the ring

Rendezvous (highest-random-weight) hashing scores every server for each key with hash(key, server) and picks the highest score. It also moves only about 1/N of keys when membership changes and needs no ring, but each lookup costs O(N).

Jump consistent hash (Google, 2014) maps a key to one of N numbered buckets with almost no memory and very even distribution, but servers can only be added or removed at the end of the list, which suits storage shards more than arbitrary failures. Load balancers such as Google's Maglev use lookup-table variants tuned for speed. The ring remains the most common choice because it handles arbitrary joins, failures and replication naturally.

Tradeoffs

DecisionUpsideDownside
Ring vs modulo placementOnly about 1/N of keys move on a membership change.Lookups need a sorted structure instead of one arithmetic operation, and plain rings can be badly unbalanced.
Decentralised placement vs controlNo lookup service on the request path.Every client must agree on membership; a stale view sends requests to the wrong server.

Real World

SystemHow it's used
Akamai (origin)Consistent hashing was first described for distributing web-cache content across changing sets of cache servers.
DiscordRoutes each guild (server) to a process using a hash ring so membership changes move few guilds.

Interview

Questions interviewers ask

  • What happens to a cache cluster using modulo hashing when you add a server?
  • Walk through a key lookup on a hash ring.
  • How many keys move when a server joins?

What a strong answer covers

Explain the remapping problem with numbers, the clockwise rule, and that about 1/N of keys move.

Common traps

  • Saying no keys move — the new server's arc does move.
  • Hashing servers by IP so every restart reshuffles keys.

Quiz

With hash(key) % N, roughly what share of keys move when N goes from 4 to 5?
  1. About 20%
  2. About 50%
  3. About 80%
  4. None

A key keeps its server only if hash % 4 equals hash % 5, which holds for about 1 in 5 keys, so roughly 80% move.

On a hash ring, which server owns a key?
  1. The closest server in either direction
  2. The first server clockwise from the key
  3. The server with the lowest load
  4. A random server

Each key belongs to its clockwise successor.

Adding one server to a ring of N moves about how many of K keys?
  1. K
  2. K/2
  3. About K/N
  4. Zero

The new server takes over only the arc before it, which on average is 1/N of the ring.

Why hash a stable node name instead of its IP?
  1. Names hash faster
  2. So restarts on a new IP keep the same position and keys
  3. IPs cannot be hashed
  4. To avoid replication

If the position depends on the IP, every reschedule moves the server around the ring and reshuffles its keys.

A server on the ring fails. Who takes its keys?
  1. All servers equally (without virtual nodes)
  2. Its clockwise successor
  3. Its counter-clockwise predecessor
  4. No one; the keys are lost

Without virtual nodes, the failed server's whole arc passes to the next server clockwise — which is exactly the imbalance virtual nodes fix.

Virtual Nodes

Give each server many positions on the ring so load spreads evenly and a failure is absorbed by many servers instead of one.

This section is part of the full PRISM roadmap, with worked examples, trade-off tables, interview questions and a quiz.

Unlock the full lesson

Rebalancing & Replication on the Ring

Store each key on several servers along the ring, and move data safely when servers join, leave or fail.

This section is part of the full PRISM roadmap, with worked examples, trade-off tables, interview questions and a quiz.

Unlock the full lesson

Practice consistent hashing in PRISM

Concepts stick when you watch them fail. Build an architecture that depends on consistent hashing, push traffic through it in the PRISM simulator, and see the latency and error rates change as you adjust the design.