Scaling Distributed Caches with McRouter: Managing the Thundering Herd and Routing Latency
Learn how to use McRouter to scale distributed Memcached clusters, prevent cache stampedes using leases, and implement consistent hashing for high-read workloads.
12 Mar 2026, 05:17 UTC

The Challenge of Massive-Scale Caching
When scaling a read-heavy application to millions of requests per second, a standard Memcached cluster often becomes a bottleneck. The primary problem is not just memory capacity, but the coordination of requests. Without a routing layer, clients must maintain a complex map of every cache server, and a single expired popular key can trigger a "cache stampede" (also known as the thundering herd problem), where thousands of concurrent requests hit the backend database simultaneously to refresh the same piece of data.
The solution is to decouple the client from the cache servers using a routing layer like McRouter. This architecture moves the logic for request distribution, failover, and consistency from the application code into a dedicated proxy layer.
Routing Strategy Comparison
Choosing how to route requests depends on whether you prioritize absolute data consistency or maximum availability. The following table compares the primary routing patterns used in distributed caching environments.
| Strategy | Mechanism | Primary Benefit | Trade-off |
|---|---|---|---|
| Consistent Hashing | Maps keys to a virtual ring of servers. | Minimizes key reshuffling during scaling. | Potential for uneven load distribution. |
| Replicated | Writes to all servers; reads from one. | High availability for critical keys. | High write latency and memory waste. |
| Pooled | Distributes requests across a pool of servers. | Efficient resource utilization. | No guarantee that a key is on a specific server. |
Solving the Cache Stampede with Leases
A critical engineering decision in this architecture is the implementation of leases. In a standard cache-aside pattern, when a key expires, the first request that notices the miss goes to the database. However, in high-concurrency environments, hundreds of requests may notice the miss at the same millisecond.
Leases prevent this by granting a "lease" to the first client that misses the cache. While that client fetches the data from the database, subsequent requests for the same key are told to wait or return a slightly stale value. This ensures that only one request hits the backend per expired key, protecting the database from crashing during traffic spikes.
Implementing McRouter Configuration
To deploy McRouter, you must define a configuration file that specifies how keys are mapped to server pools. This configuration is typically deployed on the application servers or as a separate proxy tier.
Configuration Example
# McRouter configuration snippet
# Define the physical cache servers
cache_servers = {
"pool1": {
"servers": [
"10.0.0.1:11211",
"10.0.0.2:11211",
"10.0.0.3:11211"
]
}
}
# Define the routing logic using consistent hashing
routing_table = {
"default": {
"type": "consistent_hash",
"pools": ["pool1"]
}
}
# Enable lease mechanism to prevent thundering herd
lease_enabled = true
Deployment and Verification
Run the McRouter binary on your proxy nodes with the following permissions and checks:
- Permissions: The process requires network permissions to bind to the listening port (typically 11211 or a custom proxy port) and outbound access to the Memcached pool.
- Execution: Run
mrouter --config=mrouter.conf. - Verification: Use a Memcached client (like
telnetormemcached-tool) to set a key via the McRouter port and verify it is stored on the correct backend server based on the hashing algorithm.
Limitations and Operational Risks
While McRouter simplifies client logic, it introduces specific risks:
- Memory Fragmentation: Memcached uses slab allocation, where memory is divided into fixed-size chunks. If your data sizes vary wildly, you may experience "premature eviction," where items are deleted even if total memory is available because a specific slab size is full.
- Eventual Consistency: In replicated setups, there is a window where different clients may see different versions of a key. If strict consistency is required, the application must bypass the cache for those specific operations.
- Proxy Latency: Adding a routing layer adds a network hop. This is usually offset by the efficiency of consistent hashing and reduced database load, but it must be measured in the p99 latency.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.