System designCore6 min3 things to try

Consistent hashing

How to add a server to a cluster so that only a sliver of the keys move, instead of almost all of them.

Start

1 · The problem

Add one server, and almost everything moves

The obvious way to pick a server for a key is hash(key) % N. It spreads keys evenly, until N changes. Add a server and see where the keys land now.

hash(key) % N
S1
k01k16k18k22k23
S2
k04k08k09k14k24
S3
k02k03k11k17k19
S4
k05k06k07k10k12k13k15k20k21
Keys that moved0 of 24Press “+ server”.

Every moved key is a cache miss that falls through to the database, all at once. In a cache cluster that is a thundering herd, at exactly the moment you were adding capacity.

2 · The idea

Put keys and servers on the same ring

Hash the servers and the keys onto the same circle. A key belongs to the first server clockwise from it. When a server joins, it only takes the keys between itself and the server before it.

Hash ring · 24 keys
S2S1S3
  • S1 owns 7 keys
  • S2 owns 13 keys
  • S3 owns 4 keys

Ring

0

keys moved

Modulo, same change

0

keys moved

Keys that moved get a red outline. On average only 1/N of them move.

3 · How a lookup works

Finding a key's server

  1. Hash the key

    to a position on the ring, a number between 0 and 2³².

  2. Walk clockwise

    to the next server position. Clients keep the sorted list of server positions, so this is a binary search: O(log N).

  3. Wrap around

    past the top: the first server on the ring owns the keys after the last one.

  4. Replicate

    by also writing to the next one or two distinct servers clockwise. That is how a database keeps a copy when a node dies.

4 · The catch

Few servers make lopsided arcs: use virtual nodes

With one position per server, the gaps between them are random, so one server can own twice its share. Give each server many positions on the ring and the load evens out.

4 servers · 2,000 keys
S1496
S2821
S3205
S4478
Busiest server1.64×of its fair share. The mark on each bar is a fair share: 500 keys.

5 · In a real system

Where the URL Shortener uses it

URL Shortener System Design

Cassandra places every short link on a token ring

Each Cassandra node owns ranges of the ring, with several virtual nodes each, and a link is stored on its owner plus the next replicas clockwise. Adding a node takes over a share of the ring from its neighbours; the rest of the links don't move, so the redirect path keeps working while the cluster grows.

6 · When to use it

And what else there is

ApproachGood atWatch out for
Consistent hashingCaches and stores where nodes come and goNeeds virtual nodes for even load; range queries are scattered
Modulo hashingA fixed number of serversAlmost every key moves when the count changes
Rendezvous hashingEven spread without virtual nodesEach lookup scores every server: O(N)
Range partitioningRange scans, such as everything created todayHot ranges, and a rebalancer to run

Check yourself

3 questions

1. A cluster of 5 cache servers uses hash(key) % N. You add a 6th. Roughly how many keys change server?
2. On a ring, which server owns a key?
3. Why add virtual nodes?

Takeaways

Remember this

  • Modulo hashing moves almost every key when the server count changes; a ring moves about 1/N of them.
  • A key belongs to the first server clockwise; replicas go to the next distinct servers.
  • Virtual nodes even out the load and let a bigger machine take more positions.