Distributed Cache (Memcached / Redis Cluster)
Sharpened prompt. Design a distributed in-memory cache serving 10M requests/sec across 200 nodes holding 20 TB, with P99 under 1ms, surviving node loss without a database meltdown, and supporting single keys hot enough to exceed one machine's network card.
This is the problem that appears inside every other problem on this site. If you can explain admission control, stampede prevention, and rebalancing precisely, you can hand-wave far less everywhere else.
1. Problem framing
Functional requirements
GET,SETwith TTL,DELETE, atomic counters, and compare-and-set.- Client-side routing (no proxy hop on the fast path).
- Automatic rebalancing on node join/leave.
Non-functional requirements
| Property | Target | Consequence |
|---|---|---|
| Latency | P99 < 1ms intra-AZ | One network round trip, no disk, no coordination |
| Hit rate | > 95% steady state | Admission and eviction policy is a first-class design decision |
| Availability | Cache down must not take the DB down | Load shedding and stampede control are required, not optional |
| Consistency | Best-effort; stale reads allowed | We are a cache. If you need linearizability, you need a database |
| Node loss impact | Only 1/N keys move | Consistent hashing, not modulo hashing |
Back-of-the-envelope
Working set: 20 TB across 200 nodes = 100 GB/node (r6g.4xlarge class)
Throughput: 10M ops/sec / 200 nodes = 50k ops/sec/node (Redis does ~100k+/core)
Object size: avg 2 KB -> ~50M objects/node, ~10B objects total
Network: 10M × 2 KB = 20 GB/s aggregate = 100 MB/s per node — comfortable,
unless one key is hot, which is exactly problem 5.4.
Memory index: ~50–100 bytes/object of overhead -> 5% of RAM is bookkeeping
2. High-level architecture
The client router is the interesting box. Putting routing in a smart client removes a proxy hop (and its P99 tail) from every operation; the cost is that every client must agree on topology, which is what the control plane exists to guarantee.
3. Component inventory
| Component | Concrete choice | Why this one |
|---|---|---|
| Cache engine | Redis Cluster, or Memcached for pure LRU KV | Redis when you need data structures and replication; Memcached when you need multi-threaded raw throughput and slab simplicity |
| Hash function | MurmurHash3 / xxHash (Ketama layout) | Fast, uniform, non-cryptographic. Never use Java's String.hashCode() — it clusters badly on short keys |
| Topology store | etcd with watches | Clients watch a single versioned key; topology changes propagate in under a second |
| Membership | Redis Cluster gossip, or a dedicated health checker | Gossip is self-healing; a single health checker is a SPOF |
| Near-cache | Caffeine / Ristretto | Kills hot-key traffic before it reaches the network |
| Coordination for fills | singleflight (Go) / AsyncLoadingCache (Java) | One database query per key per miss, not thousands |
| Proxy (optional) | Twemproxy / Envoy Redis filter | Only when clients are polyglot and you cannot ship a smart client everywhere |
4. The toughest parts
4.1 Rebalancing without a cascade of misses
Why it's hard. The obvious router is node = hash(key) % N. Add one node and N changes, so almost every key maps somewhere new. Every lookup misses at once, 10M requests/sec fall through to the origin database, and you have converted a capacity addition into an outage. This is the single most common cache design error.
Solution — consistent hashing with virtual nodes. Map both nodes and keys onto a 2^32 ring; a key belongs to the first node clockwise. Adding a node steals keys only from its immediate successor, so 1/N keys move instead of all of them.
Plain consistent hashing has a second-order problem: with 10 physical nodes placed randomly, ring arcs vary wildly and load can differ 3× between nodes. Fix it by giving each physical node 128–256 virtual nodes scattered around the ring — the arcs average out and load lands within a few percent of uniform.
type Ring struct {
hashes []uint32 // sorted vnode positions
owners map[uint32]string // position -> physical node
}
const vnodesPerNode = 256
func (r *Ring) Add(node string) {
for i := 0; i < vnodesPerNode; i++ {
h := murmur3.Sum32([]byte(fmt.Sprintf("%s#%d", node, i)))
r.hashes = append(r.hashes, h)
r.owners[h] = node
}
slices.Sort(r.hashes)
}
func (r *Ring) Get(key string) string {
h := murmur3.Sum32([]byte(key))
// first vnode clockwise from h; wrap at the end of the ring
i := sort.Search(len(r.hashes), func(i int) bool { return r.hashes[i] >= h })
if i == len(r.hashes) { i = 0 }
return r.owners[r.hashes[i]]
}
Worth naming as an alternative: rendezvous (highest-random-weight) hashing computes hash(key, node) for every node and picks the max. It gives perfect balance with no vnode bookkeeping at O(N) per lookup — excellent below ~100 nodes, worse above. Jump consistent hash is O(log N) and allocation-free but cannot handle arbitrary node removal, only shrinking from the tail.
The subtle failure to volunteer: during a topology change, different clients briefly hold different ring versions, so the same key is read from two nodes. For a cache this is a hit-rate blip, not a correctness bug — but only because you never write authoritative state to a cache. Say that.
4.2 Cache stampede: one expiry, fifty thousand database queries
Why it's hard. A hot key expires. In the microsecond after, every in-flight request misses, and all of them independently query the origin. The database sees a coordinated spike, slows down, which lengthens the window, which admits more requests — a positive feedback loop. TTL jitter helps with synchronised expiry across many keys but does nothing for a single hot key.
Solution — three layers, in order of preference.
(a) Request coalescing (singleflight). Only the first caller fetches; everyone else waits on the same in-flight promise. This is the highest-value 10 lines in the whole design.
var g singleflight.Group
func Get(ctx context.Context, key string) ([]byte, error) {
if v, ok := cache.Get(key); ok { return v, nil }
v, err, _ := g.Do(key, func() (any, error) { // N callers, 1 origin query
b, err := db.Load(ctx, key)
if err == nil { cache.Set(key, b, ttlWithJitter()) }
return b, err
})
if err != nil { return nil, err }
return v.([]byte), nil
}
Note that singleflight is per-process. Across 500 app instances you still get 500 origin queries, not 50,000 — usually enough. When it is not, add a distributed mutex (SET lock:key token NX PX 3000) where the loser returns stale data rather than blocking.
(b) Probabilistic early expiry (XFetch). Rather than expiring on a cliff, each reader recomputes with a probability that rises as the TTL approaches, so exactly one reader (probabilistically) refreshes before the key is gone and nobody ever sees a miss.
import random, math, time
def should_refresh(delta_ms, expiry_ts, beta=1.0):
# delta_ms = how long the last recompute took; expensive keys refresh earlier
return time.time() - delta_ms * beta * math.log(random.random()) >= expiry_ts
(c) Serve-stale-while-revalidate. Keep two clocks per entry: a soft TTL and a hard TTL. Between them, serve the stale value instantly and kick off an async refresh. The origin never sees a burst, and users never see a stall. Caffeine's refreshAfterWrite is exactly this.
4.3 Memory fragmentation and the OOM that arrives at 3am
Why it's hard. Cache values are arbitrary sizes and constantly churn. A general-purpose allocator (malloc) leaves the heap pockmarked: you have 4 GB free in total but no contiguous 8 KB block, so an allocation fails or the RSS creeps past the container limit and the OOM killer takes the node — dumping 100 GB of warm cache and stampeding the origin (see 4.2, now at cluster scale).
Solution — slab allocation. Pre-carve memory into pages, assign each page to a slab class holding fixed-size chunks (96 B, 120 B, 152 B, … growing by a factor of ~1.25), and round every value up to the nearest class. Allocation and free become O(1) pointer moves off a free list, and external fragmentation is structurally impossible.
Slab class 3 (152-byte chunks) Slab class 7 (464-byte chunks)
┌────┬────┬────┬────┬────┐ ┌──────┬──────┬──────┐
│used│free│used│used│free│ │ used │ free │ used │
└────┴────┴────┴────┴────┘ └──────┴──────┴──────┘
1 MB page, 6,898 chunks 1 MB page, 2,259 chunks
The cost is internal fragmentation: a 200-byte value in a 240-byte class wastes 40 bytes, ~17%. That is the trade, and it is a good one — bounded, predictable waste beats unbounded, unpredictable failure.
The follow-on problem is slab calcification: a workload that filled class 3 for a week and then shifts to 500-byte values finds all its memory locked in the wrong class. Memcached's automatic slab rebalancer moves pages between classes based on eviction pressure; Redis sidesteps the issue by using jemalloc with size-class bins plus activedefrag, which relocates objects during idle cycles.
What to actually say: "I'd set maxmemory to ~75% of container RAM, not 95%. Redis's own bookkeeping, replication buffers, and copy-on-write during a fork for RDB persistence can transiently double memory; the 25% headroom is what stops a background save from OOM-killing the node."
4.4 The hot key that sharding cannot fix
Why it's hard. Consistent hashing distributes keys, not load. If a single key — a viral post, a feature flag read on every request, the config blob — takes 2M requests/sec, it lives on exactly one node by definition. That node's CPU and NIC saturate while 199 others idle. Adding nodes does nothing. This is the failure mode that consistent hashing is most often wrongly believed to solve.
Solution — three tools, applied in order.
- Near-cache (the real answer). A W-TinyLFU L1 in each app process with a 1-second TTL turns 2M requests/sec into
number_of_app_instancesrequests/sec. For a 500-instance fleet, that is 500 ops/sec. The cost is up to 1 second of staleness, which for a feature flag or a view count is free. - Key replication with a salt. Write the value to
key#0…key#9, and have readers pick a random suffix. Load spreads across 10 nodes; the cost is 10× memory and a 10-way invalidation. - Detect it automatically. Sample 1 in 1000 requests into a Count-Min Sketch per client; when a key crosses a threshold, promote it into the near-cache tier dynamically. This is what Facebook's
mcrouterand Netflix's EVCache do, and naming it lands well.
4.5 Invalidation: the second hard thing
Why it's hard. The origin row changes. Now some subset of 200 cache nodes, 500 near-caches, and possibly a CDN hold a stale copy. The classic buggy pattern is update database, then update cache: two concurrent writers can interleave so the cache ends up holding the older value permanently.
Writer A: read v1 ──────────────► write cache v1 (stale, and it sticks)
Writer B: read v2 ► write cache v2 ►
Time ─────────────────────────────────────►
Solution — invalidate, don't update, and get the order right. Use cache-aside with delete: write the database first, then delete the cache key. A delete is idempotent and order-insensitive in a way that a set is not; the next reader repopulates from the source of truth. The remaining window (a reader that loaded before the write and sets after the delete) is narrow and closable with delayed double delete — delete, write DB, sleep ~500ms asynchronously, delete again.
For near-caches, publish invalidations on a fan-out channel:
# Writer
await db.update(row)
await redis.delete(f"user:{uid}")
await redis.publish("invalidate", f"user:{uid}") # every app instance subscribes
# Every app instance
async for msg in pubsub.listen():
local_l1.invalidate(msg["data"])
Redis Pub/Sub is at-most-once — a subscriber that is disconnected for 200ms misses the message and keeps stale data forever. Two fixes worth naming: bound the L1 TTL to a few seconds so any missed invalidation self-heals, or use Redis client-side caching (RESP3 tracking), where the server itself tracks which client cached which key and pushes an invalidation on change.
The strongest version of this answer: "For anything where a stale read is a correctness problem rather than a cosmetic one, I don't invalidate — I version. The cache key includes a version (user:123:v7), and a write bumps the version. Stale entries become unreachable instantly and age out on their own."
4.6 Surviving node loss without taking down the database
Why it's hard. A node holding 100 GB dies. Consistent hashing dutifully remaps its keys to the successor — which now serves double the traffic with none of it cached, misses everything, and forwards 100k requests/sec to the origin. The successor then falls over from the load, and its keys remap to its successor. This is a cascading ring collapse, and it is how cache outages become database outages.
Solution.
- Replicate the hot tier. Two replicas per shard with the client reading from either; a primary loss costs no hit rate at all. Costs 2× memory — justify it by how much origin capacity you would otherwise have to keep idle.
- Load shed at the origin, not the cache. Put an adaptive concurrency limiter (Netflix
concurrency-limits, or a simple token bucket sized to the DB's known safe QPS) in front of the origin. Excess requests get a fast error or a stale value. The database survives, degraded, instead of dying. - Circuit-break the cache client with a short timeout (5–20ms). A cache node that is slow rather than dead is more dangerous, because it consumes the app's connection pool. Fail fast to the origin.
- Warm before serving. A restarted node joins the ring only after prewarming from a peer or replaying a recent snapshot. Redis's
replicaofhandles this; joining cold is what turns a routine restart into an incident.
4.7 Eviction policy: why plain LRU is the wrong default
Why it's hard. LRU assumes recency implies importance. A scan — a crawler, a nightly batch job, a backfill — touches millions of keys once each, and LRU cheerfully admits every one of them at the head of the list, evicting the genuinely hot working set. Hit rate collapses precisely when load is highest.
Solution — admission control. W-TinyLFU (Caffeine, Ristretto, and Redis's allkeys-lfu in spirit) maintains a Count-Min Sketch of recent access frequencies and, on every insert, compares the candidate's estimated frequency against the eviction victim's. A one-hit-wonder loses and is never admitted. The sketch is tiny (4 bits per counter, halved periodically so the window slides) and typically buys 5–15 points of hit rate over LRU on real Zipfian traffic.
This has its own page, because it is a common follow-up in its own right: Cache eviction: LRU vs W-TinyLFU.
5. What breaks first
| Event | First failure | Mitigation |
|---|---|---|
| Node added during peak | Hit-rate dip on 1/N keys | vnodes bound the blast radius; warm before ring entry |
| Single viral key | One node's NIC saturates | Near-cache; key replication with salt |
| Hot key TTL expiry | Origin QPS spike | singleflight + XFetch + stale-while-revalidate |
| Workload shifts value size | Slab calcification, spurious evictions | Automatic slab rebalancing; activedefrag |
| Cache tier fully down | Origin sees 100% of traffic | Adaptive concurrency limiter at the origin, plus a degraded read path |
| Network partition | Split ring views, duplicate fills | Cache-only state means this is a hit-rate cost, never a correctness one |
6. Cheat sheet
- Routing: consistent hashing, 256 vnodes/node, MurmurHash3. Never modulo. Rendezvous hashing below ~100 nodes.
- Stampede: singleflight first, XFetch second, stale-while-revalidate third. TTL jitter for correlated expiry.
- Memory: slab classes trade bounded internal fragmentation for zero external fragmentation.
maxmemoryat 75%. - Hot keys: consistent hashing does not fix them. Near-cache does. Salted replication if you must.
- Invalidation: delete, don't update; write DB first; version keys when staleness is a correctness bug.
- Eviction: W-TinyLFU with admission control, not LRU. Scan resistance is the reason.
- The one-liner: "A distributed cache is three independent problems — placement, admission, and invalidation — and most outages come from placement solving only key distribution while load stays skewed."