Directory or Consistent Hashing
We now have two ways for a router to find a shard. Shopend’s router looks up the shop in a directory. The router for t.co computes the code’s position on the ring.
| Directory | Consistent hashing | |
|---|---|---|
| Find a key’s shard | Look it up | Compute it |
| Place a key | On any shard we choose | Where its hash lands |
| Add a shard | Move the keys we choose | Move the keys on its arcs |
| What we maintain | An entry for every key | The shards’ positions |
Shopend
For Shopend, we keep the directory. We give each national retailer a dedicated shard, and we place smaller shops on shared shards based on their workloads. With consistent hashing, a shop goes where its hash lands, so we cannot give a national retailer a shard of its own. Consistent hashing gives each shard about the same number of shops. Shopend’s shops differ in size, so the same number of shops does not mean the same workload.
For Shopend, the directory is cheap. It has one entry for each of the 10,000 shops. Each router keeps a copy of the entries it has looked up.
t.co and Drive
Consistent hashing fits when there are many keys, the keys have similar workloads, and no key needs a particular shard. The codes meet all three. A directory like Shopend’s, with one entry per shard key, would need an entry for every code, and there are billions of codes. No code needs a particular shard. Apart from viral links, which the cache handles, the codes have similar workloads. So about the same number of codes on each shard means about the same load on each shard.
For t.co and Drive, we use consistent hashing with virtual nodes instead of hash(code) mod N.
Combining hashing and a directory
We can also combine the two approaches. We hash each code into one of a fixed number of buckets, then use a directory to find which shard holds that bucket. The directory needs one entry per bucket, not per code. When we add a shard, we can move some buckets to it and update their directory entries. The codes stay in the same buckets because the number of buckets does not change.
Redis Cluster uses this approach. It hashes keys into 16,384 fixed buckets called hash slots and maintains a mapping from slots to nodes. When we add a node, we can move some slots and their keys to it and update the mapping. Each key still belongs to the same slot; only the node holding that slot changes.