On this page
Tracks

System Design — Distributed Cache

Last reviewed 11 Sept 2026

Part of the system design series. See the framework and building blocks first if you haven’t.

1. Requirements

Functional

  • Fast key-based get/set in front of a slower system of record (usually a database).
  • Support expiration (TTL) and explicit invalidation on writes.
  • Scale beyond one machine’s RAM — shard across a cluster of cache nodes.

Non-functional

  • Sub-millisecond to low-single-digit-ms latency — the entire reason it exists.
  • High hit ratio for the hot working set (the actual goal metric, not “cache everything”).
  • Must not become a single point of failure that takes the database down with it if it disappears — this is the question every interviewer eventually asks.

2. Where it sits

flowchart LR
C[Client] --> App[App server]
App -->|1: GET key| Cache[(Cache cluster)]
Cache -- hit --> App
Cache -- miss --> App
App -->|2: on miss, read| DB[(Database)]
App -->|3: SET key, TTL| Cache
App -->|write path: DELETE key, then write| DB
Cache-aside is the default read/write path

The app owns the cache-consistency logic — the cache itself is a dumb, fast key-value layer with no idea what a “stale” value means. That responsibility sits in the app’s read/write path, which is why the pattern choice below matters more than the cache technology choice.

3. Patterns and eviction

PatternRead pathWrite pathTrade-off
Cache-aside (lazy loading)App checks cache, falls back to DB, backfills cacheApp writes DB, then deletes (not updates) the cache keySimplest, most common; first read after a write is always a miss
Write-throughApp checks cache, falls back to DBApp writes cache and DB together (cache is the front door)Cache always warm; adds write latency, cache becomes part of the write’s critical path
Write-behind (write-back)Same as write-throughApp writes cache; cache asynchronously flushes to DBFastest writes; risk of data loss if the cache dies before flushing
Read-throughCache itself loads from DB on miss (app never talks to DB directly)Usually paired with write-throughCleanest app code; couples the cache implementation to the DB schema

Eviction, when the cache is full: LRU (evict least-recently-used) is the default and works for most access patterns where recency predicts future access; LFU suits a stable hot set where popularity matters more than recency; TTL is a safety net layered on top of either, so a value that should be evicted eventually is even if it stays “hot” by the eviction policy’s own metric.

4. Deep dive — invalidation races and the stampede problem

The core hard problem in caching is invalidation, not storage. Two specific failure patterns show up constantly:

Race condition on cache-aside delete-then-write. Thread A reads a miss, starts fetching from DB. Meanwhile Thread B writes a new value to the DB and deletes the cache key. Thread A’s stale DB read then lands in the cache after B’s delete, leaving a stale value cached with no further write to correct it until the TTL expires. The mitigation used in practice — the delayed double-delete pattern: delete the cache key, write the DB, then schedule a second delete of the same key ~500ms later to clear out anything a straggling read slipped in during the window.

Cache stampede (thundering herd). A very hot key expires or gets evicted, and hundreds of concurrent requests all miss simultaneously and all hammer the database to refetch the same value at once — the exact spike the cache exists to prevent. Mitigations, roughly in order of how often they show up:

  • Single-flight / request coalescing: only the first miss actually queries the DB; concurrent misses for the same key wait on that one in-flight fetch and share the result.
  • TTL jitter: randomize each key’s TTL by a few percent so a batch of keys set at the same time doesn’t all expire in the same millisecond.
  • Early recompute / probabilistic early expiration: refresh a key slightly before it actually expires, with rising probability as it nears expiry, so a background refresh usually beats the stampede.
  • Serve-stale-while-revalidate: return the (slightly) stale cached value immediately while a single background fetch refreshes it, rather than making every caller wait on a cold DB hit.

5. What real systems do today

  • Netflix EVCache — a memcached-based cache that is explicitly a Tier-0 system at Netflix (a failure directly impacts customer streaming), running across roughly 18,000 servers holding on the order of 14 petabytes of data and absorbing north of 30 million requests/sec at peak. It shards with Ketama consistent hashing and deliberately keeps all instances for one cluster within a single availability zone, replicating across AZs at the application layer rather than relying on cross-AZ hops for a single request — trading some cross-zone symmetry for tight intra-zone latency.
  • Twitter / X (historically) and many large memcached deployments — run memcached as a pure lookaside cache (cache-aside) in front of a sharded MySQL/storage tier, with client libraries handling consistent-hash routing across the memcached pool so individual node loss only reshuffles a small fraction of keys.
  • Common production guidance converging across recent engineering write-ups (2025–2026) — cache-aside plus delete-on-write plus a TTL safety net is repeatedly described as the default, boring, correct choice for most services; write-through/write-behind are reserved for cases where the cache must never be cold (e.g. a session store) because the added write-path complexity and failure coupling aren’t worth it otherwise.
  • HybridCache-style two-layer caching (in-process L1 + distributed L2 like Redis) is an increasingly common pattern for high-throughput services: an in-process cache absorbs the hottest fraction of traffic with zero network hop, falling back to the distributed cache, which falls back to the DB — reducing both cache-cluster load and the blast radius of a distributed-cache outage.

6. Scaling & failure

  • Bottleneck: single cache node’s memory/CPU ceiling → fix: shard across a cluster via consistent hashing (Ketama or similar) so growth adds nodes instead of vertically scaling one box; new cost: cross-node consistency for any operation touching multiple keys (most caches sidestep this by simply not offering multi-key transactions).
  • Bottleneck: network round trips to the cache add up on hot paths with many small lookups → fix: batch reads (MGET) instead of N sequential GETs, and/or add an in-process L1 cache in front of the distributed L2 for the hottest fraction of keys.
  • Bottleneck: a few keys get wildly disproportionate traffic → fix: this is a per-key hot-spot problem, not a capacity problem — replicate that specific key to multiple cache nodes and fan reads across the copies, or move it to a local in-process cache entirely.

Failure mode: the cache goes down — does the database fall over?

This is the question this page exists to answer directly, because it’s the one interviewers ask every time.

Yes, it can — and it’s happened to real companies. If the cache normally absorbs, say, 95% of read traffic, then a cache outage means the database instantly receives roughly 20x its normal read load. A database provisioned for the cached traffic level, not the raw traffic level, will fall over — connection pool exhaustion, disk I/O saturation, replica lag spiking, cascading timeouts back up the stack. This is functionally identical to a cache stampede, just triggered by total cache loss instead of one hot key expiring.

What actually prevents this in production:

  • Request coalescing / single-flight at the app layer, independent of the cache — so even against a cold or dead cache, N concurrent requests for the same key become 1 DB query, not N.
  • A circuit breaker in front of the DB that trips under sustained overload and serves a degraded response (stale data, a “try again” error, or a cached-at-the-CDN fallback) rather than letting every request queue up and take the DB down with it.
  • Rate limiting / load shedding at the app tier during a known cache outage — deliberately reject a fraction of requests rather than let 100% of a 20x spike reach the DB and kill it for everyone.
  • Gradual cache warm-up on recovery — bringing a cold cache cluster back at full traffic immediately just recreates the stampede in miniature for every key at once; production runbooks typically ramp traffic back or pre-warm from a snapshot/replica rather than flipping it on cold.
  • Read replicas sized with headroom for exactly this scenario — some shops explicitly provision DB read replicas assuming the cache could disappear, treating the cache as a latency optimization the DB must survive without, not a load-bearing wall the DB depends on.

Interview follow-ups

  • “Cache-aside vs write-through — when would you pick each?” — Cache-aside is the default (simple, tolerates cache loss gracefully since DB is always the source of truth); write-through when the cache must never be cold, e.g. a session store, at the cost of coupling writes to cache availability.
  • “Why delete the cache key on write instead of updating it?” — Update-on-write races with a concurrent stale read finishing after the delete; delete-then-miss-then-repopulate has no such overwrite window.
  • “What’s a cache stampede and how do you prevent it?” — Many concurrent misses for one hot key hammering the DB at once; single-flight coalescing, TTL jitter, and serve-stale-while-revalidate are the standard fixes.
  • “The cache cluster just went down entirely. What happens to your database, concretely?” — Name the multiplier (traffic the cache was absorbing lands directly on the DB), then the mitigations: coalescing, circuit breaker, load shedding, gradual warm-up on recovery.
  • “LRU vs LFU — how do you choose?” — LRU for general recency-predicts-future-access workloads (the default); LFU when there’s a genuinely stable hot set where popularity, not recency, should decide what stays.
  • “How do you scale a cache cluster horizontally without a full rehash on every node add?” — Consistent hashing (see the consistent hashing notes) — Netflix’s EVCache uses Ketama for exactly this reason.
  • “Would you ever put a cache directly in the write path (write-through) for everything, to keep it simple?” — No — it makes the cache load-bearing for every write’s latency and availability, which defeats the “cache loss should degrade, not break, the system” goal; reserve it for data that genuinely must never be cold.

Sources: Announcing EVCache: Distributed in-memory datastore for Cloud — Netflix TechBlog · Caching for a Global Netflix — Netflix TechBlog · Solving the Distributed Cache Invalidation Problem with Redis and HybridCache · System Design: Distributed Caching — Redis vs Memcached, Cache-Aside, Write-Through, Eviction, Cache Stampede · Building a Distributed Cache System: Caching Strategies and Invalidation Patterns — Medium