System design·cachingintermediate6 min

Thundering herd vs cache stampede: what breaks and five ways to fix it

A cache stampede happens when one popular cached value expires and thousands of readers miss at the same moment, so they all hit the database together. It's one case of the thundering herd problem, where many waiters wake up for one resource and all but one of them wake up for nothing.

Most read-heavy systems put a cache in front of the database for exactly this reason. Say a feed gets 10,000 reads a second and Redis answers 99.5% of them. Postgres only sees the other half percent, about 50 queries a second. So the database gets sized for 50, because 50 is what it sees on every normal day. That sizing is fine until the one value everybody reads goes missing.

Fig 01 · One hot key expires

Scroll →

Postgres queries/s

50

Cache hit rate

99.5%

p99 latency

3 ms

01
Redis answers 99.5% of reads in about 3 ms, so Postgres sees roughly 50 queries a second. Hover any box for its numbers.

The figure runs that feed. Step through it, or press Play all. In step 3, watch the queue build up beside Postgres and the hit rate fall to zero.

Thundering herd vs cache stampede#

The thundering herd is the older and more general name. It comes from operating systems. Picture eight worker processes all waiting for the same event, such as one lock being released. When it's released, the kernel wakes all eight. One of them gets the lock. The other seven find it taken and go back to sleep. Each of them paid for a wake-up and a context switch, the CPU work of swapping one process out for another, and got nothing for it.

Fig 02 · A lock release wakes every waiter

Scroll →

01

Every waiter wakes for one lock. Seven of the eight wake for nothing and pay a context switch.

Old web servers ran into this with many processes waiting to accept connections on one port. Every new connection woke all of them. Kernels fixed it by waking one waiter at a time.

A cache stampede is the same shape, moved to a cache. The resource is one cache entry. The waiters are every request that wants it. When the entry expires, they all miss at once and all go to the database for the same value.

People use the two names loosely, and in an interview either one is understood. If you want to be precise, say the stampede is the cache version of the herd.

The cache version is worse. In the lock case, the losers go back to sleep. In the cache case, every loser runs the full database query.

Why the database cannot recover on its own#

Here are the numbers from the figure. The feed gets 10,000 reads a second. Rebuilding the cached value takes a query that needs 0.8 seconds of database CPU. The database has two cores, so it can do 2 seconds of CPU work every second. Clients give up after 3.5 seconds.

Start with a normal day. Two seconds of CPU a second, divided by 0.8 seconds per query, is two and a half rebuilds a second. That's plenty, because a normal day needs at most one.

Now the key expires. In the first tenth of a second, 1,000 readers miss. Each one starts the same 0.8-second query. The database doesn't run them one after another. It shares its two cores among all of them, so each query moves at 2/1,000 of full speed. At that pace, one query needs 400 seconds to finish. The timeout is 3.5.

So no query finishes. Every client hits its timeout and gets an error. Because no query finished, nobody writes the value back to the cache. The next reader misses too, and starts another query. The database stays at 100% CPU doing work that will all be thrown away. The p99, the time the slowest one in a hundred requests takes, sits at the timeout.

It gets worse when clients retry. A reader that timed out usually tries again, so each failure adds one more query to the pile. The system won't settle by itself. It needs the traffic to drop, or someone to put the value back.

That's the part interviewers listen for. A stampede is a feedback loop. More load makes each query slower, slower queries time out, and timeouts mean the cache never refills.

Five fixes and when each one fits#

All five fixes stop many readers from rebuilding the same value at the same time. They differ in who waits, and in whether readers see an old value while the rebuild runs.

FixWhat it doesFits whenCost
Single-flightOne request rebuilds while identical ones wait for its resultReaders need fresh data and the rebuild takes about a secondWaiters sit through one rebuild
Serve stale, refresh in the backgroundKeeps the old value past its TTL and refreshes it behind the scenesA few seconds of old data is fineReaders see slightly old data
Probabilistic early refreshReads near expiry refresh early, with a chance that rises as expiry nearsVery hot keys with a known rebuild timeA few extra rebuilds and some tuning
Rebuild lock with a short expiryThe first miss sets a lock in Redis and the others wait or serve staleMany app servers share one cacheOne more Redis call per miss
TTL jitterAdds a small random amount to each TTLMany keys were written at the same momentDoes nothing for one hot key

Single-flight means one request does the work while identical ones wait for its result. In the figure's last step, the first reader to miss runs the query. The other readers attach to that same pending result. The query has the database to itself, so it finishes in 0.8 seconds. Then the value is back in Redis and everyone reads from it again.

If you have twelve app servers and each one does single-flight on its own, the database sees twelve queries. Twelve is fine.

Serving stale while you refresh is the other strong default. You keep the old value a little past its TTL, which is how long a cached value is allowed to live. The first reader after expiry starts a refresh in the background. Everyone, including that reader, gets the old value right away. Nobody waits. In exchange, readers see data that's a few seconds out of date. For a feed, a product page or a leaderboard, that's usually fine.

Probabilistic early refresh is a refinement. As a value gets close to its expiry, each read has a small chance of refreshing it early. The chance grows as expiry gets nearer. One reader almost always refreshes the value before it actually expires. You need to know roughly how long a rebuild takes to tune it.

A rebuild lock is single-flight across all your app servers. The first reader to miss sets a lock key in Redis with a short expiry, for example five seconds. Readers that find the lock either wait briefly and read again, or serve an old copy. The short expiry matters. If the server holding the lock crashes, the lock frees itself and someone else can rebuild.

TTL jitter adds a small random amount to each TTL. It's the fix for many keys expiring together, for example after a deploy warmed 100,000 entries in the same second. With jitter, their expiry spreads across a few minutes and the misses arrive as a trickle.

What to say in an interview#

When a design has a cache in front of a database, the interviewer may ask what happens when a hot entry expires. A strong answer has four parts.

  1. Name it. Say it's a cache stampede, the cache case of the thundering herd.
  2. Give the numbers. Normal load on the database is the read rate times the miss rate, say 50 queries a second. During a stampede it's the full read rate, 10,000.
  3. Explain why it doesn't recover. Queries share the database, none finishes before the timeout, and the value is never written back.
  4. Pick a fix and say why. Use single-flight if readers need fresh data, and serve stale if a few seconds of old data is fine. Add jitter if many keys were written at once.

Then say how you'd see it coming. Alert on a sudden jump in cache misses for one key, and on database CPU. A stampede shows up as both at once.

Practice this in the app: graded drills, with the misses brought back.

Start free

Related guides

Updated

Thundering herd vs cache stampede: what breaks and five ways to fix it · Archletics