A 10-node cache tier running at a 92% hit rate is saturating CPU during a traffic peak. An engineer doubles it to 20 nodes on the same consistent-hashing ring. What happens over the next two minutes?
Show the full answer Hide the answer
Second by second
Adding 10 nodes to a ring of 10 reassigns roughly half of all keys (1 − 10/20), and the new owners hold none of them. Reads for the moved half miss.
- Seconds 0–5. Aggregate hit rate falls to about 0.5 × 92% ≈ 46%. Miss rate goes from 8% to about 54%, so origin load multiplies by roughly 6.8× against an origin sized for 8% of traffic.
- Seconds 5–20. The origin saturates. Its latency rises, so each miss holds an application thread for longer, so the application's concurrency limit fills and its own queue grows.
- Seconds 20–60. Request timeouts fire. If clients retry, offered load on the origin rises again while the cache is still mostly empty, because the cache can only refill as fast as the origin can serve, and the origin is now the bottleneck.
- Where it amplifies: without per-key coalescing, every concurrent caller for the same missing hot key issues its own origin query. The hottest keys generate the most duplicate work at exactly the worst moment.
What the user sees: latency and errors rising sharply, starting the moment capacity was added.
Why the refill is slower than people expect
Refill time ≈ working-set bytes ÷ (origin throughput × mean item size). A 40 GB working set behind an origin that can supply 20,000 items/second at 4 KB is roughly 40 GB ÷ 80 MB/s ≈ 8 minutes of degradation, and that is the optimistic case where the origin never falls over.
What stops it
A mechanism, not vigilance:
- Single-flight per key so N concurrent misses become one origin query and N−1 waiters.
- Serve stale while revalidating. A moved key that has a stale copy anywhere is a latency problem, not an origin problem.
- Admission control at the origin, so it degrades at its limit instead of collapsing past it.
- Warm new nodes before they join the ring — replay recent keys into them, then add them.
- Add nodes in small increments. Ten nodes going to eleven moves about 9% of keys, not 50%.
The structural answer
For a stateful tier whose value is the data it already holds, scale vertically first. Larger nodes on the same ring add memory, CPU and network without moving a single key, and a 10-node tier has plenty of headroom before instance size becomes the constraint. Scale out when one node's memory or network caps you, and then do it before the peak with warming, never during it.
The general rule: horizontal scaling is cheap for stateless tiers and expensive for stateful ones, because the cost is paid in data movement rather than in instances.
What would have to be true for it to self-heal
The origin would have to absorb 6.8× its design load for the whole refill window without its own latency rising. That is a legitimate design — an origin provisioned for a cold cache is the strongest form of static stability — and it is expensive, which is why almost nobody has it and almost everybody assumes they do.
Common weak answers
- "More cache nodes means more cache, so it gets better." True at steady state, after the rebalance, which is minutes away. The question was about the next two minutes.
- "Autoscale the origin too." Instance start-up plus warm-up is minutes, and the surge is seconds.
- "Switch the cache to a bigger instance type later." That is the right answer, offered at the wrong time: the move itself is a restart, which empties the cache.