Consistent Hashing

In this chapter, we change how t.co computes a link’s shard, so that adding a shard moves only the links the new shard takes over.

After reading this chapter, you should be able to:

  • Explain why adding a shard with hash(code) mod N moves most links, and estimate how many move
  • Place shards and keys on a hash ring, and find the shard for a key
  • Define consistent hashing
  • Determine which keys move when a shard is added to the ring or removed from it
  • Explain why one position per shard can leave the shards with uneven shares of the keys
  • Define virtual nodes, and explain how they spread the keys a shard gains or gives up across several shards
  • Give a larger server a larger share of the keys
  • Compare a directory with consistent hashing, and explain why Shopend keeps its directory

Sections

  1. Adding a Shard
  2. The Hash Ring
  3. Changing the Shards on the Ring
  4. Virtual Nodes
  5. Directory or Consistent Hashing