Skip to content

Consistent hashing

Mapping keys to nodes so that adding or removing a node moves only a small share of keys, instead of reshuffling almost all of them.

Distribution

Learn it

0 of 1 checks done
  1. The obvious way to spread keys over N servers is hash(key) % N. It balances well, until N changes. Going from 10 to 11 cache servers changes the server for most keys, so the hit rate collapses and every key is fetched from the database at once.

    Anything that routes by key (caches, connection servers, shards) needs a mapping that survives membership changes.

  2. Consistent hashing places nodes and keys on the same circular hash space, a ring. A key belongs to the first node clockwise from its hash. Adding a node takes over only the keys between it and its predecessor; removing one hands its keys to the next. On average only 1/N of keys move.

  3. Add a node under each placement, and compare one point per node with a hundred.

    2,000 keys spread over cache nodes. Add a node and count the keys that now belong somewhere else: each one is a cache miss, or data to copy.
    Placement
    4 nodes
    4 nodeskeys by hash mod N

    Keys per node. The line is a perfectly even share.

    • Node 1500
    • Node 2508
    • Node 3512
    • Node 4480

    The busiest node holds 1.0× an even share.

  4. Check

    What do virtual nodes (many ring points per physical node) improve?

Quick reference

The same ideas, condensed for revision.

How it goes wrong

Hot key
One key's load exceeds a node; hashing cannot split it.
Disagreeing membership
Clients with different views of the ring send the same key to different nodes.
Too few virtual nodes
Uneven arcs give some nodes several times the load of others.

Instead, consider

Modulo hashing
The number of nodes never changes, or reshuffling everything is cheap.
Directory / range map
You need to move specific ranges deliberately (e.g. to rebalance hot shards) rather than by hash.
Rendezvous hashing
Node counts are small and you want simple, even placement without virtual nodes.

In practice

Cassandra / ScyllaDB / DynamoDB
Token rings with virtual nodes and replication.
Client-side memcache routing
Consistent hashing in the client library or a router such as mcrouter.
Load balancer hash policies
Route by a header or path so one key's requests reach one backend.

It assumes

  • Clients (or a router) agree on the current membership of the ring.
  • Keys are numerous and individually small compared with a node's capacity.

Explain it in your own words

Write at least 60 characters (0 so far). Write it as you would say it in a design review. You will compare it against the points a strong answer makes.

Where you practise it

Further reading

Engineers describing it in systems they run.

  • Partitioning

    Splitting data or work by key so each part is handled independently: scaling out, and giving each key a single owner.

  • Caching

    Keeping a copy of data closer to where it is used, trading freshness and complexity for speed and reduced load on the source.

  • Replication

    Keeping copies of data on several machines for durability, read capacity and locality, and living with copies that briefly disagree.

  • Request coalescing

    When many callers ask for the same thing at the same moment, do the expensive work once and give every caller the result. The fix for thundering herds and hot keys.