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
Select a component to see what it is responsible for and which state it owns.
- 1Users → Web servers: Page request
- 2Web servers → mcrouter: get / multiget, delete
- 3mcrouter → memcached pool: Keys by consistent hash
- 4Web servers → MySQL: Query on miss; writes
What you need to know
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.
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. 4 nodesKeys 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.
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.
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?