Distributed Caches Explained: Consistent Hashing, Eviction & the Thundering Herd
Ten thousand reads a second against a database that answers in ten milliseconds — the arithmetic of why caches exist, and how a ring, some replicas, and three eviction policies turn one box of RAM into a planetary system.
Explain it like I’m five
Imagine a big school where hundreds of kids keep asking the librarian for the same five books. Walking to the main library every time is slow, so each classroom keeps a little shelf with the most-asked-for books right by the door.
That shelf is a cache. Now make it tricky: there are thirty classrooms, and you need rules. Which classroom keeps which book? Give every book a number, arrange the classrooms in a circle, and each book lives in the first classroom clockwise from its number — that is consistent hashing. Keep a spare copy in the next classroom too, in case one room is locked — that is replication. When a shelf fills up, put back the book nobody touched for weeks — that is eviction. And when the popular book’s loan slip expires, do not let all thirty kids stampede the library at once — that is the thundering herd, and this guide is the set of rules that stops it.
Intuition: the question that separates senior from staff
“Design a distributed cache” shows up at exactly the interviews where “I’d add Redis” is not an answer — Meta, Google, Amazon, Apple. The interviewer is not testing whether you know a cache exists. They are testing whether you can reason about placement, failure, full memory, and expiry stampedes, in that order. By the end of this guide you will derive all four, and know how the real systems build them.
Start with the arithmetic that makes the cache inevitable. Your database answers a read in 10 milliseconds. Traffic is 10,000 reads per second. Little’s law says the database must sustain 10,000 × 0.01 = 100 concurrent reads — one hundred open connections doing nothing but reads, before a single write. Now put a cache in front with a 95% hit rate: only 500 reads per second reach the database, and the concurrent load drops to 5. The cache does not just make reads fast; it divides your database load by twenty.
And the workload is on your side. Reads are wildly skewed: roughly 80% of reads ask for the same 1% of data. Cache the hot 1% and you capture most of the win; the other 99% barely matters. This skew is why a small box of RAM can front a huge database — and why every design decision below is about managing the hot keys.
Four questions structure the whole design, and they arrive in interview order. Placement: with many cache boxes, which box owns which key? Failure: a box dies at 2am — what breaks? Full memory: RAM is finite — what gets thrown away? Expiry: a hot key’s TTL fires — what stops a thousand requests from stampeding the database at once? Answer those four and you have designed a distributed cache.
Spread the keys with consistent hashing so adding a node moves only 1/N of them, replicate every key three ways so a dead node is a non-event, evict with LRU plus TTL, and coalesce the thundering herd so one request rebuilds while the rest wait.
How the four pieces work
Consistent hashing: the ring
One cache box is a bottleneck and a single point of failure, so you run N of them — and now every key needs an owner. The naive rule is hash(key) mod N. It works until the fleet changes: add a ninth box and N changes, so nearly every key remaps to a different box. Your cache effectively empties itself on every scaling event — a full fleet of cold misses slamming the database at once.
The ring fixes it. Hash both the keys and the nodes onto a circle. Each key belongs to the first node clockwise from its position. Add a node and it only steals the keys in its own arc — about 1/N of the keyspace. Remove one and only its arc’s keys move, to the next node clockwise. The rest of the fleet never notices.
This is the design Amazon’s Dynamo paper made canonical: DeCandia et al. partition data with consistent hashing, where each key’s coordinator is the first node clockwise and the next N−1 nodes on the ring form the key’s preference list of replicas. In practice each physical node claims many points on the ring — virtual nodes — so arcs stay even and a beefy box can simply claim more of them.
Source: DeCandia et al., “Dynamo: Amazon’s Highly Available Key-Value Store,” SOSP 2007 ↗The ring: keys and nodes share one circle; ownership is “first node clockwise.” Topology changes disturb only the arcs they touch.
Replication: plan for failure
At fleet scale, failure is not an emergency — it is the normal case. So every key lives three times: on its home node plus the next two nodes clockwise. One node dies at 2am and the other two copies answer; no data is lost, no pager goes off, and the dead node’s arc is absorbed by its clockwise neighbor until a replacement joins the ring.
Writes go to all three replicas; reads need only one. The cost is straightforward: three times the RAM for the guarantee that a single dead box is a non-event. This is the same preference-list idea as the ring section — placement and replication are one mechanism, not two.
Eviction: when memory fills up
RAM is finite, so the cache needs a rule for what gets thrown away. Three classic policies: LRU tosses the least recently used key, LFU tosses the least frequently used, and TTL stamps an expiry on every key so stale data dies on its own. Most production caches run LRU plus TTL together — recency handles the skew, expiry bounds the staleness.
Redis makes the choice one config line, maxmemory-policy: allkeys-lru, allkeys-lfu, volatile-lru and volatile-ttl (which only evict keys that have a TTL), and noeviction. Worth knowing cold: noeviction is the default — a fresh Redis used as a cache does not evict at all; when memory fills, writes fail with an out-of-memory error. For a cache, set allkeys-lru; Redis’s own docs call it the good default when a small subset of keys gets most of the traffic.
Source: Redis docs — Key eviction ↗The thundering herd — and request coalescing
Here is the trap the interview is really about. A hot key carries 1,000 reads per second and its TTL expires. In one instant, a thousand requests miss the cache and slam the database with the same query — the spike the cache existed to prevent, arriving all at once. The database that was comfortably serving 500 reads per second is suddenly asked for 1,000 identical ones.
The fix is request coalescing: the first request through does the database fetch while the other 999 wait on its result. One query, not a thousand. The alternative is serve-stale: keep serving the expired value while a background refresh rebuilds it — slightly stale answers, zero spike. Facebook’s Memcached paper describes the same idea as leases: only the client holding the lease performs the database fetch, and everyone else waits rather than stampeding.
A hot key expiring is a thousand requests becoming a thousand database queries. Coalesce them: one rebuilds, the rest wait — or serve the stale value and refresh behind it.
Consistency: delete, don’t update
On a write, do you update the cached copy or delete it? Delete it. Write the new value to the database, delete the cache key, and let the next read re-fetch and re-populate. Updating the cached copy in place risks a race where a stale write lands after a fresh one; invalidation has no such race — a deleted key can never serve a lie. This pattern is called cache-aside, and it is the default answer unless the interviewer asks for something stronger.
Cold start: warm it before it takes traffic
A new node joins the ring knowing nothing — every key in its arc misses until it fills. If it takes full traffic immediately, those misses become a mini thundering herd. So warm it first: pre-load the hottest keys (the top 1% you already know about), or let it shadow traffic and fill gradually before it serves real reads. The interviewer asks this right after you say “add a node, only 1/N keys move” — the moved keys still have to be fetched from somewhere.
Two systems that live on this
This is not textbook theory. The largest cache fleets on earth run exactly these four mechanisms — here is how two of them map the ideas onto production reality.
Fleet: Facebook’s Memcached
Facebook’s NSDI 2013 paper “Scaling Memcache at Facebook” describes a fleet serving billions of requests per second over trillions of items. Keys are distributed across the fleet by consistent hashing; clients talk to any server in an all-to-all pattern, so no proxy tier can bottleneck. For failure, roughly 1% of the fleet runs as a Gutter pool that takes over for sick nodes. For consistency, invalidation daemons watch the databases and delete cached copies on write — delete, don’t update, at planetary scale. And for the thundering herd, clients use leases: only the lease-holder does the database fetch while everyone else waits.
Source: Nishtala et al., “Scaling Memcache at Facebook,” NSDI 2013 ↗Edge: Cloudflare’s tiered cache
Cloudflare puts cache in hundreds of cities close to visitors — and then caches the caches. With Tiered Cache, lower-tier data centers near the visitor serve repeat requests, while a smaller set of upper-tier hubs aggregate the misses; only the upper tiers ever ask your origin. A hot asset fetched once in Frankfurt serves the whole region instead of every city hitting your server. For freshness, Cloudflare supports stale-while-revalidate: serve the cached copy immediately while a background fetch refreshes it — the serve-stale answer to the thundering herd, built into the edge.
Source: Cloudflare docs — Tiered Cache ↗The eviction dial: Redis
Redis is the cache most teams actually operate, and its eviction behavior is one config line away from a production incident. The maxmemory-policy setting picks the policy — allkeys-lru, allkeys-lfu, volatile-ttl — and the default is noeviction, which refuses writes when memory fills instead of evicting. Every “Redis is down” story that is really “Redis is full” traces back to this line. Set allkeys-lru when you deploy Redis as a cache, and say so in the interview before they ask.
Source: Redis docs — Key eviction ↗Work the cache by hand
Interviewers trust candidates who do the arithmetic out loud. Here is the full cache design reduced to five lines of math you can reproduce on a whiteboard.
Deploying Redis as a cache with default settings. Memory fills, maxmemory-policy is noeviction, and writes start failing with out-of-memory errors instead of evicting old keys. The fix is one line — allkeys-lru — and naming it unprompted is the kind of detail that ends interviews early.
The ring, the eviction, and the singleflight — in Python
Three mechanisms, thirty lines each. The ring decides placement, the OrderedDict decides eviction, and the in-flight map decides that one request rebuilds while the rest wait.
import hashlib
from collections import OrderedDict
from threading import Lock, Event
class ConsistentHashRing:
"""Keys walk the ring to the first node clockwise."""
def __init__(self, nodes, replicas=100):
self.ring = {}
for node in nodes:
for i in range(replicas): # virtual nodes
self.ring[self._hash(f"{node}:{i}")] = node
self.sorted = sorted(self.ring)
@staticmethod
def _hash(key):
return int(hashlib.md5(key.encode()).hexdigest(), 16)
def node_for(self, key):
h = self._hash(key)
for point in self.sorted:
if point >= h:
return self.ring[point]
return self.ring[self.sorted[0]] # wrap around
class LRUCache:
"""OrderedDict: front is hot, back is evicted."""
def __init__(self, capacity):
self.store = OrderedDict()
self.capacity = capacity
def get(self, key):
if key not in self.store:
return None
self.store.move_to_end(key, last=False)
return self.store[key]
def put(self, key, value):
self.store[key] = value
self.store.move_to_end(key, last=False)
if len(self.store) > self.capacity:
self.store.popitem(last=True) # evict coldest
class CoalescingCache:
"""One rebuild per key; everyone else waits on it."""
def __init__(self, capacity):
self.cache = LRUCache(capacity)
self.inflight = {}
self.lock = Lock()
def get(self, key, loader):
hit = self.cache.get(key)
if hit is not None:
return hit
with self.lock:
if key in self.inflight:
event = self.inflight[key]
owner = False
else:
event = Event()
self.inflight[key] = event
owner = True
if owner:
value = loader(key) # the one DB query
self.cache.put(key, value)
with self.lock:
del self.inflight[key]
event.set()
return value
event.wait() # the 999 wait here
return self.cache.get(key)The pattern has a name in Go’s standard library — singleflight — and the idea ports anywhere: a per-key in-flight marker turns a thousand duplicate queries into one.
How this gets asked
Six prompts drawn from how Meta, Google, Amazon, and Apple actually probe this topic — each with the answer shape, the follow-up, and the trap.
Start with the math: 10,000 reads/sec at 10 ms is 100 concurrent database connections — that is the problem the cache exists to erase. Then walk the four questions in order. Placement: consistent hashing ring, first node clockwise, so scaling moves only 1/N keys. Failure: replicate every key 3x on the next two clockwise nodes, so a dead node is a non-event. Full memory: LRU plus TTL, and if it is Redis, set maxmemory-policy to allkeys-lru because noeviction is the default. Expiry: request coalescing so one fetch rebuilds the hot key while the rest wait. Close with cache-aside: write to the DB, delete the key, let the next read re-populate.
“What if the interviewer says the workload is write-heavy?” — Answer: then the cache buys little and invalidation dominates; say you would shorten TTLs and consider write-through for the hot keys.
Name it first: a thundering herd. One thousand reads per second all miss at the same instant and become one thousand identical database queries — the exact spike the cache was built to absorb. The fix is request coalescing: the first request through performs the fetch while the others block on its result, so the database sees one query. Mention the alternative — serve the stale value and refresh in the background — and the precedent: Facebook’s Memcached paper calls the same mechanism leases, where only the lease-holder hits the database.
“What if the rebuild itself is slow?” — Answer: then waiting requests pile up; bound the wait with a timeout and fall back to the stale value.
“Nothing breaks, and here is why.” Every key is stored three times — its home node plus the next two clockwise — so the dead node’s arc is still fully readable from the replicas. Reads for those keys keep hitting the surviving copies; the next node clockwise absorbs the dead node’s range. What does change: the fleet is down one copy of every key in that arc, so a second failure in the same arc would hurt — which is why a replacement node joins the ring and the data re-replicates before morning. Mention the Dynamo framing: treat failure as the normal case, not the emergency.
“Two nodes in the same arc die?” — Answer: the keys whose three replicas were exactly those nodes are unavailable until one returns; that is the calculated risk of N=3, and you would raise it for data you cannot lose.
Delete it. Write the new value to the database, delete the cache key, and let the next read re-fetch and re-populate — the cache-aside pattern. Updating in place has a race: a stale write can land after a fresh one and the cache serves the wrong value indefinitely. A deleted key can never serve a lie; the worst case is one cache miss. Add the nuance: if the key is extremely hot, you accept the miss storm risk by coalescing the re-fetch — the herd fix from earlier composes with this answer.
“When would you update in place?” — Answer: write-through, when reads must never miss — but you pay the write latency on every write and you own the race.
Start with the workload: LRU assumes recency predicts reuse, LFU assumes frequency does. For skewed read traffic — the 80/20 kind caches are built for — both work, and LRU wins on simplicity: an OrderedDict, O(1), no counters to maintain. LFU earns its keep when the hot set is stable but recency lies — a nightly batch job that touches a million cold keys once would flush an LRU cache and evict the actually-hot keys; LFU’s frequency counters survive that. Then the practical close: either way, pair it with TTLs for a staleness bound, and in Redis set maxmemory-policy explicitly because the default noeviction evicts nothing.
“How do you pick the TTL?” — Answer: from the business cost of staleness — seconds for prices, minutes for profiles, longer for content that rarely changes.
About one-ninth. On the consistent hashing ring the new node only claims the keys in its own arc — the keys for which it becomes the first node clockwise. The other eight nodes’ arcs are untouched. Contrast with hash-mod-N, where changing N remaps nearly every key and the fleet cold-starts at once. Then volunteer the part they were waiting for: the moved keys still have to be fetched from somewhere, so the new node warms its arc — pre-load the known-hot keys — before it takes full traffic, or the “only 1/N moved” win turns into a mini thundering herd.
“What if the nodes have different sizes?” — Answer: virtual nodes — the big box claims more points on the ring and proportionally larger arcs.
Takeaways
- Do the read math first. 10,000 reads/sec × 10 ms = 100 concurrent connections — the cache divides database load, not just latency.
- The ring owns placement. First node clockwise; topology changes move only 1/N keys. Virtual nodes keep arcs even.
- Replicate 3x and sleep. Home node plus the next two clockwise — a dead node at 2am is a non-event.
- Set your eviction policy on purpose. LRU + TTL is the default answer; Redis’s default noeviction evicts nothing — change it.
- Coalesce the herd. One request rebuilds the expired hot key; the rest wait. Or serve stale and refresh behind.
- Delete, don’t update. Cache-aside invalidation can never serve a lie; a stale in-place update can.