Skip to content

Design a Distributed Cache (Memcache), stage 5 of 9: break it

A cache server dies

The database is provisioned for the normal miss rate. Decide what clients do in the minutes before the replacement arrives.

System so far· 5 parts
1234CLIENTUsersSERVICEWeb serversSERVICEmcrouterCACHEmemcached poolDATABASEMySQL

Select a component to see what it is responsible for and which state it owns.

  1. 1Users → Web servers: Page request
  2. 2Web servers → mcrouter: get / multiget, delete
  3. 3mcrouter → memcached pool: Keys by consistent hash
  4. 4Web servers → MySQL: Query on miss; writes

What you need to know

0 of 1 checks done
  1. Clients choose a cache server for each key by hashing the key. Two common schemes:

    • hash mod N: server = hash(key) % N. Simple, but changing N moves almost every key.
    • Consistent hashing: servers and keys are placed on a ring; each key belongs to the next server clockwise. Adding or removing a server moves only the keys next to it. See Consistent hashing.
  2. Add a node under each placement and compare how many keys move. Then try the ring with one point per node versus 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.

  3. Moving few keys is good for planned changes. A failure raises a different question: the dead server's keys all miss, and their load must go somewhere. Rehashing sends them to the neighbouring servers, which are already busy. And keys aren't equal: one hot key can be a fifth of a server's traffic.

  4. Check

    A dead server's hottest key (20% of its traffic) is rehashed onto a healthy server that is already at 85% capacity. What can happen?