The Hash Ring
With hash(code) mod N, the number of shards is part of the formula. When that number changes, the formula gives a different shard for most codes.
If we take mod N out, we are left with hash(code). That number does not depend on the number of shards. But it is a hash value, not a shard. We still need a way to map each hash value to a shard.
Consistent hashing
In 1997, David Karger and his colleagues at MIT were working on spreading cached web pages across a group of cache servers. They were using hashing to assign pages to servers, and they had the same problem as we do with shards. Adding or removing one server would move most pages to a different server. Each moved page would then miss in the cache on its new server.
Their solution was consistent hashing. With consistent hashing, a key stays on its server when a server is added or removed. The only keys that move are the ones the new server takes over, or the ones the removed server held.
Codes on a ring
A hash function returns a number in a fixed range. A real one might return a number from 0 to . That is about 4 billion. The hash values in our examples, in this chapter and the previous one, are made up, and each of them has at most two digits. So let’s assume our hash function returns a number from 0 to 99.
Without mod N, we use this whole range. Picture the numbers from 0 to 99 placed around a circle, the way the minutes are placed around a clock. Going clockwise, 99 is followed by 0. We call this circle the hash ring.
Now we place the t.co short codes on the ring. Each code goes on the ring at its hash value. Here are the five codes from the previous section, with their made-up hash values:
| Code | hash(code) |
|---|---|
x7Kp2Qa |
17 |
b3Rt9Lm |
52 |
Qz81fWc |
30 |
Hn4vR8e |
71 |
pL2wJ5s |
96 |
Shards on the same ring
Assume for a moment that we had 100 shards, placed at positions 0, 1, 2, …, 99. Then each code would go to the shard at the code’s position.
In reality, the range of hash values is much larger than the number of shards. So most positions on the ring have no shard. In our example, we place three shards on the 100 positions.
To assign a code to a shard, we take these steps:
-
We put the shards on the same ring. To find a shard’s position, we hash the shard’s name with the same hash function. Suppose A lands at 20, B at 55, and C at 85.
-
Each code goes to the first shard we reach by moving clockwise from the code’s position.
x7Kp2Qais at 17. Moving clockwise, the first shard we reach is A at 20. So the link is on A.pL2wJ5sis at 96. Moving clockwise, we pass 99, continue from 0, and reach A at 20. So that link is on A too.
Here is our example of five codes and three shards:
| Code or shard name | Position | Stored on |
|---|---|---|
x7Kp2Qa |
17 | A |
| Shard A | 20 | |
Qz81fWc |
30 | B |
b3Rt9Lm |
52 | B |
| Shard B | 55 | |
Hn4vR8e |
71 | C |
| Shard C | 85 | |
pL2wJ5s |
96 | A |
Each shard holds the codes on the arc that ends at its position. A holds the codes at positions 86 to 99 and 0 to 20. B holds the codes at positions 21 to 55. C holds the codes at positions 56 to 85.
Finding the shard
The router keeps the shards’ positions in a sorted list: A at 20, B at 55, C at 85. To route a request for a code, the router computes the code’s hash. Then it finds the first position in the list that is greater than or equal to the hash. If there is no such position, it takes the first position in the list. That is the same as moving clockwise past 99 and continuing from 0.
The list has one entry per shard, so it is small. If we have multiple shard routers, each router keeps a copy. When we add or remove a shard, we update every router’s copy.