Skip to main content
Create your own
Lesson illustration

Preventing Cache Stampede: Locking & Probabilistic Recomputation

Hello! Welcome back to our module on Advanced Distributed Concepts.

In our last lesson, we focused on generating unique identifiers at scale using the Snowflake algorithm, ensuring every entity in our distributed system can be uniquely named without conflict. Today, we shift our focus from creating new things to accessing existing ones—specifically, what happens when a resource becomes too popular and our safeguards (like caching) suddenly fail.

This brings us to the learning outcome for this lesson: to diagnose and mitigate cache stampede (thundering herd) using techniques like locking or probabilistic recomputation. This is a classic failure mode in distributed systems, and understanding how to prevent it is a hallmark of a seasoned engineer. It's also a frequent topic in system design interviews, making it a crucial addition to your toolkit.

The Problem: When Your Shield Becomes a Weapon

In most systems, a cache is introduced to improve performance and protect the database from excessive load. It acts as a fast, in-memory shield. But what happens when that shield suddenly vanishes?

A cache stampede, also known as a thundering herd, occurs when a very popular cached item (a "hot key") expires. At that moment, all the concurrent requests that were happily being served by the cache suddenly result in a cache miss. In a high-traffic system, this could mean hundreds or thousands of application threads simultaneously trying to regenerate the data, typically by querying the database.

This diagram illustrates a typical setup where a cache stampede can occur. Multiple application servers, all serving requests for the same popular data, find the cache key has expired. They all proceed to query the database simultaneously, overwhelming its connection pool and processing capacity.

This situation is particularly insidious because the very component you relied on for protection—the cache—becomes the trigger for a system-wide failure.

One of the most important mental shifts for a systems engineer is to see the cache not just as a performance enhancer, but as a critical, load-bearing part of the architecture. Let's explore this idea further.

uRadical blog - Four Concepts That Separate Systems Thinkers From Syntax Chasers

This article from uRadical.io offers a powerful perspective on this problem. It argues that a cache is not a mere optimization but a fundamental structural component.

Please read the section titled "02 — The Thundering Herd". Pay close attention to the author's distinction between a "performance optimisation" and a "load-bearing wall." This reframing is key to understanding the severity of the problem.

As the article highlights, this is a correlation problem. Your system might be perfectly capable of handling thousands of requests per second for different keys, but a concentrated burst of requests for the same key at the same instant can be catastrophic. The database, provisioned for a steady cache-miss rate, is suddenly exposed to the full, raw request rate, leading to connection pool exhaustion, high CPU load, and cascading failures.

Mitigation Strategies: Breaking the Synchronization

The key to preventing a thundering herd is to break the synchronization of the requests hitting the database. We need to introduce coordination or randomness to ensure that only one—or very few—requests attempt to regenerate the cache at any given time. Let's explore the most common techniques.

1. Locking (Coordination)

The most direct approach is to use a distributed lock. The logic is simple: only the process that obtains the lock is allowed to go to the database.

Caching Pitfalls Every Developer Should Know

The ByteByteGo channel provides a clear, animated explanation of the problem and the locking solution.

First, watch the introduction to Cache Stampede to see a visual representation of the problem. Then, continue with the explanation of the locking strategy, which outlines the core logic and the different options available to requests that fail to acquire the lock.

As the video explains, when a process experiences a cache miss:

  1. It attempts to acquire a distributed lock (e.g., using a SET key value NX EX command in Redis, which is atomic).
  2. If the lock is acquired: The process is now the "leader." It queries the database, repopulates the cache, and then releases the lock.
  3. If the lock is not acquired: Another process is already regenerating the cache. This process should not hit the database. Instead, it can:
    • Wait a short time and then re-check the cache (polling).
    • Serve a slightly old, or "stale," version of the data if one is available. This is the basis of the popular stale-while-revalidate pattern, which prioritizes availability.
    • Immediately return an error or a "not found" response.

This "lease" mechanism, as it's sometimes called, is a powerful form of coordination.

How to Build Cache Stampede Prevention - OneUptime

This article provides a practical, code-level view of implementing a distributed lock strategy with a stale data fallback.

Read the section on Locking Strategies. While the code is in Node.js, the logic is universal. Focus on the flow: check cache, attempt lock, on success fetch and update, on failure check for stale data or wait. Pay special attention to the "Lock Considerations," as these are crucial for a robust implementation (e.g., ensuring the lock timeout exceeds computation time).

2. Probabilistic Early Recomputation

Locking is a reactive strategy—it kicks in after the cache has expired. An alternative is to be proactive and refresh the cache before it expires. Probabilistic early recomputation does this in a clever, decentralized way.

Here's the core idea: As a cached item gets closer to its expiration time, any process reading it has a small, but increasing, probability of deciding to "do a good deed" and refresh the data. The probability is low enough that only one or two processes out of thousands will actually trigger the refresh, avoiding a stampede.

The formula can be as simple as (current_time - last_refresh_time) / TTL multiplied by a small constant. The randomness ensures that refreshes are staggered.

This sequence diagram shows the logic of probabilistic early expiration. Upon a cache hit, if the data is not expired but is nearing expiration, the application performs a probabilistic check. If the check passes, it proactively refreshes the data from the database. If it fails, it simply serves the existing (but still valid) cached data.

This approach is excellent for very hot keys with predictable traffic patterns, as it smooths out the recomputation load over time.

3. Request Coalescing (Single-Flight)

This technique is particularly relevant given your background in Go. The golang.org/x/sync/singleflight package provides a beautiful, idiomatic solution to the thundering herd problem at the process level.

Request coalescing, or single-flight, ensures that for many concurrent function calls with the same key, only the first call executes the function. All other calls wait for the first one to complete and then share its result.

uRadical blog - Four Concepts That Separate Systems Thinkers From Syntax Chasers

The same uRadical article provides a concise explanation of Facebook's "lease" mechanism and shows how Go's singleflight package implements the same core idea.

First, read about the real-world incident at Facebook and their lease-based solution. Then, read the section "The Fix Is Coordination, Not Capacity", which demonstrates the elegant singleflight implementation in Go.

As you can see from the Go example, singleflight.Group acts as a gatekeeper. When 10,000 goroutines call GetUser("some-hot-key") simultaneously, group.Do ensures that db.QueryUser is called exactly once. The other 9,999 goroutines block until that one query returns, and then they all receive the same result. This effectively prevents the thundering herd within a single application instance.

For a fully distributed system, you could combine singleflight (for in-process deduplication) with a distributed lock (for cross-process coordination) for a belt-and-braces approach.

4. Jitter: The Simple Power of Randomness

Sometimes the simplest solutions are the most effective. Jitter is the practice of adding a small amount of randomness to fixed timers.

Instead of setting a TTL of exactly 60 seconds on all related cache keys, you could set it to 60s + random(0s to 5s). This small variation desynchronizes the expiration times, causing cache misses to be spread out naturally over time instead of all occurring at the same instant.

This same principle is famously used to prevent thundering herds in retry mechanisms, as a real-world case from PayPal demonstrates.

How PayPal Beat the Thundering Herd Problem and Fixed Their Architecture

This video by Arpit Bhayani explains how PayPal solved a cascading failure problem caused by synchronized retries. The principle is identical to preventing synchronized cache expirations.

Watch the section that identifies the root cause of the thundering herd problem. Then, continue through the explanation of jitter as the solution. Notice how a seemingly complex distributed systems problem boils down to simply adding a small, random delay.

Conclusion

In this lesson, we've dissected the cache stampede problem, a critical failure mode that arises from the synchronized failure of a system's caching layer. By understanding that a cache is a load-bearing component, not just an optimization, you're better equipped to design resilient systems.

Key Takeaways:

  • Cache Stampede (Thundering Herd): A phenomenon where the expiration of a popular cache key causes a massive, simultaneous flood of requests to the backend database, overwhelming it.
  • The Cause: The problem is one of correlation. The system is brought down by many clients demanding the same resource at the same time.
  • The Goal of Mitigation: Break the synchronization of requests through coordination or randomness.
  • Key Techniques:
    • Locking: A reactive strategy where only one process gets a "lease" to regenerate the cache.
    • Probabilistic Early Recomputation: A proactive strategy that staggers cache refreshes before expiration.
    • Request Coalescing (singleflight): An in-process mechanism to deduplicate concurrent function calls for the same resource.
    • Jitter: Adding randomness to TTLs or timers to desynchronize events.

In our next lesson, we will continue exploring caching pitfalls by examining another classic problem: cache penetration. We will learn how to diagnose and mitigate this issue, where attackers or bugs cause requests for non-existent data to bypass the cache and hammer your database.

Can't find a good explanation? Sign up and we'll make it for you

Sign up