Limited time: AI code review, hints, mock interviews, whiteboard analysis, and all Pro features are unlocked. Enroll
⏱️ 62 min read

Designing a Distributed Cache

Difficulty: Intermediate Topics: Consistent Hashing, Eviction, Replication, Hot Keys Asked at: Amazon, Google, Meta, Uber Prerequisites:Caching and Consistent Hashing


1. Understanding the Problem

A distributed cache is a pool of machines that hold hot data in RAM so your database does not have to answer the same question ten thousand times a second. You SET a key with a TTL, you GET it back in under a millisecond, and when memory fills up the cache throws away whatever it thinks you need least.

The word that matters here is volatile. A cache is allowed to lose your data. It will lose your data - on eviction, on a node crash, on a TTL firing, on a failover that drops the last few writes. Everything in this design follows from accepting that. If you need the data to survive, you are designing a key-value store instead, and the answer looks completely different: write-ahead logs, quorum writes, Merkle-tree repair. None of that appears here.

Real examples: Redis Cluster, Memcached (Facebook runs it in front of their MySQL tier), Netflix EVCache, Amazon ElastiCache, Hazelcast.


2. Naive First Cut

flowchart LR
    APP["App Server"]:::client
    CACHE["Redis<br/>single node"]:::service
    DB[("Origin DB<br/>Postgres")]:::data

    APP -->|"1. GET key"| CACHE
    APP -->|"2. On miss query DB"| DB
    APP -->|"3. SET key"| CACHE

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Color Meaning
🟣 Purple Clients
🟒 Green Services
🟑 Yellow Data stores
πŸ”΅ Blue Edge / CDN

One Redis box. Every app server talks to it. Cache miss falls through to Postgres and the app writes the value back.

Why this breaks:

The rest of the doc evolves this into a shardable, survivable cache tier.


3. Prior Art We’re Drawing From


4. Functional Requirements

Core (Top 3)

  1. GET and SET with a TTL - read a key in single-digit milliseconds at p99, write a key with an optional expiry after which it disappears on its own
  2. Scale past one machine’s memory - the working set must spread across many nodes, and adding a node must not invalidate the cache
  3. Survive losing a node - one machine dying costs you a slice of the cache, not the whole thing, and the cluster keeps serving in the meantime

Below the Line


5. Non-Functional Requirements

Core

NFR Target
Read latency p99 under 5 ms end to end from the app, p50 under 1 ms
Throughput 1M ops/sec sustained at a 10:1 read-to-write split
Availability 99.99% for the tier. Losing one node degrades hit rate, never availability.
Hit rate 95%+ steady state, because the origin DB is sized for 5% of 1M ops/sec and no more

Below the Line


6. Scale Estimation (Back-of-Envelope)


7. Core Entities


8. API / System Interface

The cache protocol itself is not HTTP - it is a binary or RESP protocol on a long-lived TCP connection, because an HTTP round trip costs more than the whole operation budget. Written as an interface:

GET key
  Response: value bytes, or MISS

SET key value [EX ttl_seconds]
  Response: OK

DEL key
  Response: 1 if the key existed, 0 otherwise

MGET key1 key2 key3
  Response: list of values with MISS holes
  Note: all keys must hash to slots on one node, or the client fans out

INCRBY key delta
  Response: new integer value
  Server-side atomicity so counters do not need read-modify-write

Two operational endpoints the client library depends on:

CLUSTER SLOTS
  Response: slot ranges mapped to primary and replica addresses

MOVED <slot> <host:port>
  Not a request - an error reply. Means "that slot lives elsewhere now,
  refresh your topology and retry there."

Security note: A cache has no business being reachable from the internet. Put it in a private subnet, require AUTH, and terminate TLS between app and cache if they cross a trust boundary. The classic breach here is an unauthenticated Redis bound to 0.0.0.0 - and because the cache holds session tokens and user profiles, the blast radius is not β€œsomeone warmed my cache.”


9. High-Level Design

Build it in three steps: make one node fast, make many nodes addressable, make the loss of one node survivable.

FR1: GET and SET with TTL - One Node, In Memory

Start with the question an interviewer actually wants answered: why keep anything in RAM at all when Postgres already has a buffer pool?

Latency, by three orders of magnitude at each step:

Where the data is Rough access time
L1 or L2 CPU cache ~1 ns
Main memory (RAM) ~100 ns
NVMe SSD read ~100 Β΅s
Postgres query over the network ~1-10 ms

A DB query is not slow because disks are slow. It is slow because you pay a network round trip, connection acquisition, parse, plan, buffer-pool lookup, row assembly, and serialisation. The cache skips all of it and answers out of a hash table in the same process that owns the socket. The honest number is that an in-memory cache turns a 1 ms operation into a 0.2 ms operation, and the 0.2 ms is mostly network.

New components:

  1. Cache Server - a single-threaded event loop over a hash table. One thread is a feature, not a limitation: no lock contention, no cache-line ping-pong between cores, and every operation is atomic by construction.
    πŸ’‘ Single-threaded means operations are serialised. A slow command - KEYS * on 250M keys - blocks every other client on that node. This is the number one way people take down a Redis in production.
  2. Eviction structure - a hash table gets you O(1) lookup but tells you nothing about what to throw away when memory fills. The textbook answer is a doubly linked list threaded through the entries: on every access, unlink the node and move it to the head. The tail is then the least recently used entry, and eviction is popping the tail in O(1).
  3. Client Library - lives in the app process. Owns the connection pool, serialises values, applies timeouts, and decides what to do on a miss.
flowchart LR
    APP["App Server<br/>client library"]:::client
    CACHE["Cache Server<br/>hash map plus LRU list"]:::service
    DB[("Origin DB<br/>Postgres")]:::data

    APP -->|"1. GET key"| CACHE
    CACHE -->|"2. MISS"| APP
    APP -->|"3. Read through on miss"| DB
    APP -->|"4. SET key with TTL"| CACHE

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

Step-by-step flow:

  1. App calls cache.get("user:1234:profile") through the client library
  2. Library picks a pooled connection and writes the command. Round trip inside one AZ is ~0.3 ms.
  3. Cache server hashes the key into its table. Hit: it moves the entry to the head of the LRU list and writes the value back.
  4. Miss: the library gets a MISS and queries Postgres (~2 ms)
  5. App calls cache.set("user:1234:profile", value, ttl=300)
  6. Cache server inserts into the hash table, links the entry at the head of the LRU list, and records an expiry timestamp

How does TTL actually expire? Not with a timer per key - 250M timers is absurd. Two mechanisms together. Lazy expiry: on every read, compare the stored expiry to now, and if it has passed, delete the entry and report a miss. That makes expiry correct from the caller’s point of view the instant it fires. Active expiry: a background loop samples a few keys from the expires table every 100 ms and deletes the dead ones, so keys nobody ever reads again still get reclaimed. Lazy expiry gives correctness, active expiry gives you the memory back.


FR2: Scale Past One Machine - Sharding Without Nuking the Cache

500 GB does not fit in one box, so split the keyspace. The obvious split is the wrong one, and the reason it is wrong is the best five minutes of this whole design.

Start with node = hash(key) % N. Ten nodes, keys land evenly, reads are fast. Then you add an eleventh node because traffic grew.

Every key’s ownership is now computed mod 11 instead of mod 10. A key only stays put if hash(key) % 10 == hash(key) % 11, which happens for roughly 1 in 11 keys. About 90% of 250M keys are now looked up on a node that has never seen them. Not relocated - the data is still sitting on the old node, perfectly intact, just unreachable because nobody asks that node for it any more.

What that does to the system, in order:

  1. Hit rate falls from 95% to about 10% in the time it takes clients to refresh their topology - call it seconds.
  2. Miss traffic goes from 50K/sec to ~900K/sec, all of it aimed at an origin DB provisioned for 50K/sec.
  3. Postgres connection pool saturates, query latency climbs, app threads block waiting on the pool.
  4. App request timeouts fire. Clients retry. Retries add load.
  5. The DB is now the outage, and it stays down until you either shed traffic or warm the cache back up.

Adding capacity caused a total outage. That is the failure mode you have to design out.

The fix is consistent hashing.
πŸ’‘ Hash both keys and nodes onto the same circular space, usually 0 to 2Β³Β². A key belongs to the first node you meet walking clockwise from the key’s position. Add a node and it takes over only the arc between itself and its predecessor - about 1 of N keys move, and every other key keeps its home.

Adding an 11th node to a 10-node ring now relocates roughly 9% of keys instead of 90%. Hit rate dips from 95% to about 86%, miss traffic roughly triples rather than rising eighteen-fold, and the origin absorbs it.

Plain consistent hashing has one real flaw: with only 10 node positions on a ring of 4 billion, the arcs are wildly uneven. Hash three nodes onto the ring and one of them routinely ends up owning 50% of the space while another owns 15%. The fix is virtual nodes - give each physical node 100-200 positions instead of one.

Why 100-200 and not 5 or 5000? It is a variance argument. Each vnode’s arc length is a random variable; a physical node’s share is the sum of its vnodes’ arcs. Summing more independent samples shrinks the relative spread of the total, roughly as 1 over the square root of the count. At 1 vnode per node the worst node might carry 2-3Γ— the average. At 16 vnodes you are down to maybe Β±30%. At 150 you are inside Β±10%, which is close enough that normal traffic skew dominates anyway. Past a few hundred you are paying for ring-lookup cost and topology size to buy evenness you cannot measure.

Redis Cluster takes a cleaner variant of the same idea: 16384 fixed slots. slot = CRC16(key) mod 16384, and slots are explicitly assigned to nodes. The slot number never changes for a given key - only the slot-to-node assignment moves. 16384 slots over 16 nodes is 1024 slots each, fine-grained enough that migration happens in small pieces, and the whole slot map compresses into a bitmap small enough to gossip around constantly.

New components:

  1. Topology-aware client - caches the slot map locally and routes each key to the owning node directly, no proxy hop. On a MOVED reply it refreshes the map and retries.
  2. Shards - 16 primaries, each owning ~1024 slots and ~31 GB of the working set
flowchart LR
    APP["App Server<br/>client library"]:::client
    TOPO["Slot Map<br/>16384 slots to nodes"]:::service
    N1["Primary 1<br/>slots 0 to 1023"]:::data
    N2["Primary 2<br/>slots 1024 to 2047"]:::data
    N3["Primary 3<br/>slots 2048 to 3071"]:::data
    DB[("Origin DB")]:::data

    APP -->|"1. Hash key to slot"| TOPO
    TOPO -->|"2. Route to owner"| N1
    TOPO -->|"3. Route to owner"| N2
    TOPO -->|"4. Route to owner"| N3
    APP -->|"5. Fill on miss"| DB

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

Step-by-step flow:

  1. App calls get("user:1234:profile")
  2. Client computes CRC16(key) mod 16384 β†’ slot 9642
  3. The cached slot map resolves 9642 β†’ Primary 10 at 10.0.3.14:6379
  4. GET goes straight to that node. One network hop, no router tier.
  5. Node answers from RAM
  6. If the client’s map is stale, the node replies MOVED 9642 10.0.3.19:6379. The client refreshes CLUSTER SLOTS and retries. One wasted round trip, no error surfaced to the app.

Why no proxy tier? A proxy would add a hop, and a hop is ~0.3 ms on a 5 ms budget. Pushing routing into the client costs you a topology-refresh mechanism and fatter clients, and buys you the lowest achievable latency. Proxies (Twemproxy, Envoy with a Redis filter) earn their place when you cannot upgrade every client, or when connection count becomes the problem: 500 app servers Γ— 16 nodes Γ— a pool of 20 is 160K connections, which a proxy tier collapses.


FR3: Survive a Node Loss - Replication and Failover

One primary dies. It held 1/16 of the working set, so your hit rate drops from 95% to about 89% and ~56K extra requests/sec land on the origin. Survivable. What is not survivable is that the slots it owned now have no owner - every GET and SET for those 1024 slots returns a cluster error, not a miss, and a cluster error is not something the app’s cache-aside path knows how to absorb.

So each shard gets a replica.

Replication is asynchronous, and it has to be. A SET on the primary acknowledges the client immediately and streams to the replica in the background. Waiting for the replica would add a round trip to every write, roughly doubling write latency, to protect data you have already agreed is disposable. Nobody makes that trade for a cache. (Redis offers WAIT for semi-synchronous acknowledgement. In a cache you leave it alone.)

What a failover loses: everything in the replication backlog at the moment the primary died, typically single-digit to low-hundreds of milliseconds of writes. In practice: a handful of keys that were SET just before the crash are either absent or hold the previous value on the promoted replica. The next read misses, falls through to the origin, and refills. The cost of a failover is a small burst of misses - not corruption, as long as nothing treats the cache as authoritative.

New components:

  1. Replica nodes - one per primary, 16 of them. They hold a full copy of their primary’s slots and serve no traffic by default.
    πŸ’‘ Reading from replicas looks like free capacity and is usually a trap. The replica lags the primary by the replication delay, so a client that writes then immediately reads can see its own write vanish. Only route reads to replicas for data where you genuinely do not care.
  2. Cluster coordinator (gossip) - every node pings every other node over a dedicated cluster bus. When enough primaries independently mark a node as failed, the cluster agrees it is gone and promotes one of its replicas.
flowchart LR
    APP["App Server"]:::client
    P1["Primary 1<br/>slots 0 to 1023"]:::data
    R1["Replica 1"]:::data
    P2["Primary 2<br/>slots 1024 to 2047"]:::data
    R2["Replica 2"]:::data
    GOSSIP["Cluster Bus<br/>gossip and failure detection"]:::async

    APP -->|"1. Read and write"| P1
    APP -->|"2. Read and write"| P2
    P1 -->|"3. Async replicate"| R1
    P2 -->|"4. Async replicate"| R2
    GOSSIP -->|"5. Ping all nodes"| P1
    GOSSIP -->|"6. Promote on failure"| R1

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
    classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0

Step-by-step flow when Primary 1 dies:

  1. Primary 1 stops answering pings. Peers that notice mark it PFAIL - possibly failed, a local suspicion only.
  2. Suspicions propagate over the cluster bus. Once a majority of primaries report PFAIL for the same node, it is promoted to FAIL cluster-wide. Requiring a majority is what stops one node with a bad NIC from evicting healthy peers.
  3. Replica 1 notices its primary is marked FAIL and requests votes to take over
  4. A majority of primaries grant the vote. Replica 1 promotes itself and claims slots 0-1023.
  5. The new topology gossips out. Clients still pointed at the old address get a MOVED and refresh.
  6. Total time to recover: roughly the failure-detection window (a few seconds, tunable) plus the election. During that window, reads for those slots error out.

What the app must do during those seconds: treat a cache error exactly like a miss, with a circuit breaker so you do not burn the request budget retrying a dead node. A cache error is a cache miss with extra latency - any app that returns a 500 because Redis blinked has made the cache a hard dependency, which defeats the point of having one. Related reading: database replication covers the sync-versus-async trade in the durable case, where the answer is different.


10. Technology Choices

Tier What it stores Access pattern Primary pick Alternatives
Cache tier 500 GB of hot values with TTLs 1M ops/sec single-key GET and SET Redis Cluster Memcached, Hazelcast, Aerospike, ElastiCache
Near cache A few thousand hottest keys, per app process Up to 100K ops/sec in-process, no network Caffeine Guava Cache, Ehcache, a plain bounded map
Cluster membership Slot map, node health, epoch Low volume, read constantly by clients Redis Cluster gossip ZooKeeper, etcd, Consul
Metrics Hit rate, p99 latency, evictions, memory, per-key hotness Scrape every 15 s, query on alert Prometheus with redis_exporter Datadog, CloudWatch, VictoriaMetrics

Why Redis Cluster over Memcached here. Memcached is genuinely better at some things. It is multithreaded, so one box scales with cores instead of needing one process per core. Its slab allocator keeps per-key overhead lower, which at 250M keys is real memory. And it has no replication, no failover, no cluster bus - less code to misbehave.

Redis Cluster wins on this requirement set for three specific reasons. First, FR3 asks for node-loss survival, and Redis has replication plus automatic failover built in; with Memcached you build that yourself or accept that a dead node means a cold slice. Second, the hot-key and stampede deep dives both lean on server-side atomicity - SET NX for single-flight, INCRBY for counters, Lua for compare-and-set on a versioned key - and Memcached’s string-only surface makes those awkward. Third, slot indirection gives online resharding without the client-side ring rebuild that memcached clients handle inconsistently across languages.

Pick Memcached when the workload is pure opaque-blob caching, you have a lot of small keys and memory is the binding cost, and you can tolerate losing a shard cold. That is close to Facebook’s and Netflix’s situation, and both of them chose it deliberately. For a general service cache, Redis Cluster’s extra machinery pays for itself.

Why a near cache at all, given the cache is already sub-millisecond? Because it removes the network, and the network is the budget. A Caffeine lookup is ~100 ns against ~300 Β΅s for a Redis round trip - three orders of magnitude - and it is the only answer to a hot key that does not involve more cache nodes. The cost is staleness bounded by the near cache’s TTL, which is why you keep that TTL to a few seconds.

Why gossip over ZooKeeper for membership? ZooKeeper gives you a strongly consistent view of who is alive, and it is the right answer when a split-brain costs you correctness. Here it does not: two nodes both believing they own slot 9642 costs you duplicate work and some staleness, not a corrupted ledger. Gossip buys you one less system to operate and no external quorum in the failure path. Reach for etcd or ZooKeeper when the cache is also doing distributed locking, where a split-brain genuinely does break something.


11. Data Modeling

A cache has no schema. What it has is a keyspace convention, and the convention is load-bearing - it is the only tool you have for finding, invalidating, and accounting for entries after the fact.

Key naming:

user:1234:profile          β†’ serialised user profile, TTL 300s
user:1234:permissions      β†’ permission set, TTL 60s
product:98765:v7           β†’ product detail, version-suffixed, TTL 600s
session:a1b2c3d4           β†’ session blob, TTL 1800s
ratelimit:ip:203.0.113.7   β†’ integer counter, TTL 60s
feed:1234:page:0           β†’ precomputed first page, TTL 30s

Why namespace with colons. The scheme buys you four concrete things:

  1. Ownership is legible. user:* belongs to the profile service. When memory is tight you can answer β€œwho is using 200 GB” instead of guessing.
  2. Invalidation is targetable. A user updating their profile needs user:1234:profile gone and user:1234:permissions left alone. Flat keys like profile_1234 make that a string-munging exercise.
  3. Metrics group naturally. Tag hits and misses by prefix and you get a hit rate per entity type. A global 95% hit rate hiding a 40% hit rate on permissions:* is the kind of thing this surfaces.
  4. Collisions stop happening. Two teams both caching β€œ1234” is a genuine production incident, and it is unambiguous once prefixed.

Include a version or schema marker when the serialised shape can change - product:98765:v7. A deploy that changes the struct then reads v8, misses cleanly, and refills. Without it you deserialise v7 bytes with v8 code and get an exception on a cache hit, which is a miserable thing to debug.

Memory accounting per entry. This is where cache capacity planning actually happens, and the overhead is larger than people expect:

Component Approximate bytes
Key string, plus its length header and allocator rounding key length + ~16
Value object header and length metadata ~16
Hash table entry - key pointer, value pointer, next pointer ~24
Expiry entry in the expires table, for keys with a TTL ~24
Allocator rounding on the value 0 to 24

Call it 50-90 bytes of overhead per key before the value, and treat 70 as the planning number. The consequence is sharp: caching 250M entries of 2 KB each costs ~17 GB of pure bookkeeping, about 3% waste. Caching 250M entries of 20 bytes each - a user-id-to-flag map - costs the same 17 GB on 5 GB of data, 350% waste. Small values are where caches get expensive. The fix is to stop storing one key per item: pack the map into a single hash field-per-item structure, where the per-field overhead is a fraction of a full key’s, or batch items into one serialised blob per bucket.

Slot ownership. The mapping from key to machine runs through two levels of indirection, both deliberate:

key          "user:1234:profile"
  ↓          CRC16 of the key, mod 16384
slot         9642                     fixed forever for this key
  ↓          slot map, gossiped by the cluster
node         10.0.3.14:6379           changes on resharding or failover
  ↓
shard        primary 10.0.3.14 plus replica 10.0.4.14

16384 slots across 16 primaries is 1024 slots and ~31 GB per primary. Migration moves whole slots, so the smallest unit of rebalancing is ~31 MB - small enough to move without a visible blip.

One sharp edge: a multi-key operation only works if every key lands in the same slot. MGET user:1234:profile user:5678:profile hits two different slots and is rejected or fanned out by the client. If you need keys co-located, hash tags force it: Redis hashes only the substring inside {...}, so {user:1234}:profile and {user:1234}:permissions share a slot and can be read in one round trip. Use it sparingly - every hash tag is a hand-placed hot spot that resharding cannot break up.

Access patterns:

Operation Path How
Read a hot key Near cache β†’ primary Caffeine lookup first, then one hop to the owning primary
Read a cold key Primary β†’ origin DB Miss, single-flight fill from Postgres, SET with TTL
Write through an update Origin DB β†’ primary Commit to Postgres, then DEL the key, short TTL as backstop
Invalidate a user’s entries Primary DEL the known key names. Never KEYS user:1234:* - it blocks the event loop.
Count something Primary INCRBY server-side, so no read-modify-write race
Reshard Slot map Migrate slots one at a time, redirecting in-flight requests with ASK

12. Deep Dives

Deep Dive 1: Eviction - Why Nobody Runs True LRU

Problem: Memory is full and a SET arrives. Something has to go. Choosing badly means evicting an entry you are about to need, and every bad eviction is a future origin query. The policy is also on the hot path, so whatever you pick runs 91K times a second.

Bad: FIFO - evict the oldest entry by insertion time. It is a queue, it is trivial, and it ignores the only signal that matters. Your most popular key, read 50K times a second, gets evicted the moment it ages out, while a key written once and never read again survives because it arrived later. On a workload with any locality, FIFO’s hit rate trails LRU by a wide margin. The one place FIFO is defensible is when every entry has identical value and a hard freshness bound, and that is rare.

Good: True LRU - hash map plus a doubly linked list, move to head on access, evict the tail. Correct, O(1), and the textbook answer. Two costs show up at scale:

Great: Approximate it. Two production approaches, and they optimise for different workloads.

Sampled LRU, which is what Redis does. Instead of maintaining global ordering, store a last-access clock per entry (a few bytes, no pointers) and on eviction sample a handful of random keys, evicting the oldest of the sample. The default maxmemory-samples is 5, and from Redis 3.0 the algorithm also keeps a pool of good candidates across evictions, which pulls it much closer to true LRU than naive sampling. Redis’s own documentation is explicit that the reason it does not implement true LRU is memory cost, and that at a sample size of 10 the approximation is very close to theoretical LRU. You give up exactness - a sampled evictor will occasionally discard something true LRU would have kept - and you get the 16 bytes per entry back, no list mutation on reads, and no contention point.

W-TinyLFU, which is what Caffeine does. Shift the question from β€œwhat do I evict” to β€œdoes this new entry deserve to be admitted at all.” Maintain a compact frequency sketch - a count-min sketch, a few bits per counter - over recently seen keys. On a miss, compare the incoming key’s estimated frequency against the eviction candidate’s. Admit the newcomer only if it looks more popular. The W is a small LRU window in front, so a genuinely new hot key is not locked out by a sketch that has never seen it.

When LFU beats LRU: scan resistance. Picture a cache full of hot product pages, and then a batch job walks the entire product table once. Under LRU, every scanned row is a recent access, so the scan marches through the cache evicting your hot set in favour of entries that will never be read again. You come out the other side with a cold cache and no warning. Frequency-based admission is immune: each scanned key has a frequency of 1 and loses every admission contest, so the scan streams through without displacing anything. Any workload mixing interactive traffic with batch or analytical passes over the same keyspace wants LFU-flavoured admission. Redis exposes this as allkeys-lfu, implemented with a Morris counter per key plus a decay so yesterday’s hot key does not stay privileged forever.

LRU still wins where recency genuinely predicts reuse and the workload shifts - session data, a news feed where today’s article is hot and last week’s is dead. LFU’s decay handles that, but less directly than LRU does.

In simple terms: Keeping a perfectly ordered list of what you used last costs more than it is worth. Redis instead peeks at five random keys and throws out the stalest one, which is nearly as good for a fraction of the cost. Caffeine goes further and asks whether a new item is even popular enough to deserve a slot - so a one-off scan of your whole database cannot wash your hot data out of the cache.

Why sampling beats exactness: The point of eviction is to maximise hit rate, not to produce a correct ordering. Sampled LRU and true LRU differ by a fraction of a percent in hit rate on real power-law traffic, and that fraction costs 4 GB of pointers plus a write on every read. You spend the 4 GB on more cached entries instead, which buys more hit rate than the perfect ordering ever would.


Deep Dive 2: Hot Keys - When One Key Gets 100K Ops per Second

Problem: Consistent hashing distributes keys evenly. It says nothing about traffic. A celebrity posts, a product goes viral, a config key gets read on every request - and now one key is taking 100K ops/sec. Every one of those lands on the single primary owning that key’s slot. That node is at 62K ops/sec from normal traffic already, so it is now being asked for 162K ops/sec against a ceiling near 100K. Its event loop queue grows, p99 for that node goes from 0.5 ms to 80 ms, and crucially every other key on that node suffers too - all 1023 other slots are collateral damage. Meanwhile the other 15 primaries sit at 40% CPU. You cannot fix this by adding nodes: resharding moves the slot to a different machine, where it does exactly the same thing.

Bad: Ignore it. The usual reasoning is that hot keys are rare and the cache is fast. Both are true and neither helps, because hot keys are not random - they are caused by the product working. The launch, the viral post, the front-page placement all create them, which means the hot key arrives precisely when you can least afford a degraded cache tier. The specific way this ends: the hot node’s latency breaches the app’s cache timeout, the app treats timeouts as misses, and the miss traffic for everything on that node hits the origin.

Good: Put a local cache in front, in the app process. Caffeine with a 10K-entry cap and a 1-5 second TTL. The hot key is by definition in every app server’s local cache within milliseconds of going hot, so the 100K ops/sec collapses to one refresh per app server per TTL window. With 200 app servers and a 2-second TTL, the cache tier sees 100 ops/sec for that key instead of 100K. A thousand-fold reduction from thirty lines of code.

The limit is staleness, and it is a hard floor: a write is invisible for up to the near cache’s TTL, on every app server independently, and there is no way to invalidate it. You deleted the key in Redis; 200 local caches do not know and will not know until their copies age out. Set the TTL to 2 seconds and you have accepted up to 2 seconds of serving the old value. For a product price that is usually fine. For a permission check after a revoke, it is a security bug - which is why user:*:permissions should carry a 1-second near-cache TTL or none at all.

The second limit is memory: 200 app servers each holding 10K entries is 2M duplicated entries. Fine for the hot tail, hopeless as a general strategy.

Great: Detect hot keys and replicate them across shards, with single-flight on the fill.

Detection. You cannot replicate what you cannot see. Two places to measure:

Use both, and aggregate client reports centrally. A key crossing a threshold - say 5K ops/sec, or 2% of a node’s capacity - gets flagged.

Replication. Once flagged, write the key under N suffixed names spread across different slots:

product:98765            original, one slot, one node
product:98765#0          replica 0, different slot
product:98765#1          replica 1, different slot
...
product:98765#9          replica 9, different slot

Readers pick a suffix at random: key + "#" + random(0, 10). Ten copies on (very likely) ten different nodes means 10K ops/sec each instead of 100K on one. The suffixes land on different slots because the hash of product:98765#3 is unrelated to the hash of product:98765#7 - which is the whole trick, and the reason you must not wrap the key in a hash tag here.

Writers must now update all N copies, and they will not all land at the same instant. The staleness this buys: during a write, different readers hit different replicas and see different values - some old, some new - for as long as the fan-out takes, typically a few milliseconds. Worse, if one of the N writes fails and you do not retry, that replica stays wrong until its TTL expires. So: keep TTLs on hot replicas short (30-60 s) so a lost write self-heals, and accept that a hot key is eventually consistent with itself. For a view counter or a product page, fine. For anything where two users comparing screens matters, do not use this.

Request coalescing rides along. When the hot key does miss - TTL expiry, failover, deploy - you do not want N replicas Γ— 200 app servers all querying the origin. The client library keeps an in-process map of in-flight fetches keyed by cache key: the first thread to miss starts the load, every subsequent thread for that key attaches to the same future and waits. Go’s singleflight, Java’s CompletableFuture deduplication, Caffeine’s LoadingCache all do this natively. Per process it turns a hundred concurrent misses into one origin query; across processes you still have 200 queries, which is where Deep Dive 3 picks up.

In simple terms: One wildly popular key can melt one cache server while the rest idle. First, keep a copy inside each app server so most reads never leave the process. Second, store the popular key under ten slightly different names so the load spreads across ten servers, and read from a random one. The price is that for a few milliseconds after a write, different users may see different versions of that key.

Why replication beats just adding nodes: Adding nodes reduces the per-node share of distinct keys. A hot key is one key - it cannot be split by any hashing scheme, because hashing is deterministic and that is the entire point. The only ways to spread load on a single logical key are to change the key name (replication) or stop asking the network (near cache). Capacity is not one of the options.


Deep Dive 3: Cache Stampede - The Mechanism, Not the Headline

Problem: A key serving 50K reads/sec expires. In the microsecond after expiry, every in-flight request for it misses. All of them query the origin for the same row, concurrently. The origin was sized for 50K ops/sec across the whole system and just received 50K concurrent queries for one row. Connection pool exhausts, queries queue, latency climbs, app threads block on pool acquisition, and the app stops serving requests that have nothing to do with this key.

The caching concept page covers the shape of this and the three standard answers. What follows is the mechanism underneath each one, because the mechanism is where the interview goes.

Bad: Nothing. Plain cache-aside with a fixed TTL and no coordination. Worth knowing why this is usually invisible in testing: at low traffic, the window between β€œkey expired” and β€œfirst fetcher has written the new value” is shorter than the gap between requests, so the stampede never forms. The herd needs arrival_rate Γ— fill_latency > 1 concurrent requests to exist. At 50K reads/sec and a 20 ms origin fetch, that is 1000 concurrent requests. At 5 reads/sec it is 0.1, so you see one extra query and conclude the problem is theoretical. It is not theoretical, it is traffic-dependent, and it appears the day you get popular.

Good: Single-flight. One request does the work, the rest wait.

In-process it is a map from key to in-flight future, and that is enough for one JVM. Across 200 app servers you need a shared lock, which means a lock in the cache itself:

1. GET key                       β†’ MISS
2. SET lock:key owner_id NX EX 5 β†’ did I get it?
3. Got it      β†’ query origin, SET key value EX 300, DEL lock:key
4. Did not get β†’ sleep 20ms, GET key, retry up to a bound

The NX makes step 2 atomic - exactly one of the thousand contenders gets OK. EX 5 is the part people omit and then regret: without an expiry, a fetcher that crashes between acquiring the lock and deleting it leaves the key permanently un-fillable, and every reader spins forever. The 5 seconds must exceed your p99 origin fetch time or you get two fetchers; it must stay short or a crash costs you that long.

Three real limits:

Facebook’s memcache tier uses a sharper variant: a lease. On a miss, the server hands exactly one client a token and tells the others to wait or serve stale. One mechanism covers both the stampede and the stale-set race from Deep Dive 4, because a stale set can be rejected if its lease was invalidated by an intervening write.

Great: Never let the expiry happen all at once. Two techniques, and you want both.

Probabilistic early recomputation (XFetch). Instead of refreshing exactly at expiry, each reader independently decides whether to refresh early, with a probability that climbs as expiry approaches. Store the time the last fill took (delta) alongside the value, and on every read evaluate:

now - delta * beta * log(random(0,1)) >= expiry  β†’ refresh now

log(random(0,1)) is negative, so the term pushes the comparison earlier the larger delta is; beta (default 1.0) tunes eagerness. The behaviour that falls out: far from expiry almost nobody refreshes, near expiry the probability rises smoothly, and the expected refresh happens shortly before expiry - so the key is replaced while it is still valid and no reader ever sees a miss. Expensive keys (large delta) start refreshing earlier, which is exactly backwards from intuition and exactly right, because expensive keys are the ones you cannot afford a herd on. The paper is Optimal Probabilistic Cache Stampede Prevention (VLDB 2015); the mechanism is a dozen lines.

Plain jitter - TTL = 300 + random(0, 30) - is the poor cousin and is still worth doing. It does not prevent a herd on one key, it prevents the correlated herd where ten thousand keys cached during the same deploy all expire in the same second. That version of the problem is more common and jitter kills it outright.

Stale-while-revalidate. Decouple β€œcan I still serve this” from β€œshould I refresh this.” Store two timestamps:

value        the cached bytes
fresh_until  now + 300s   serve without thinking
stale_until  now + 600s   serve, and trigger a background refresh

Between fresh_until and stale_until, a reader returns the stale value immediately and kicks off an async refresh (guarded by single-flight so only one fires). Nobody waits. Past stale_until the entry is genuinely gone and readers block on a real fill. The win is that the origin’s latency leaves the user-facing path entirely during normal refreshes - and if the origin is down, you keep serving the stale value for the full stale_until window instead of failing. That property makes this the single most valuable technique on this page for availability, not just for load. It is the same contract as HTTP’s stale-while-revalidate directive, applied inside your own cache.

In simple terms: When a popular key expires, thousands of requests notice at the same instant and all ask the database the same question. The fix in two parts: let only one of them actually ask (the rest wait for its answer), and better still, refresh the key a moment before it expires - choosing the moment randomly, so each reader picks a slightly different one - so it never expires under load at all. And keep the old copy around a little longer than its official lifetime, so you can hand it out while the fresh one loads.

Why early recomputation beats locking: Locking is a reaction - the miss has already happened and you are limiting the damage, while every waiter pays the full fetch latency. Early recomputation is prevention: the key is replaced before it expires, so there is no miss to coordinate around and no waiter to make slow. Locking is the backstop for the cases prevention misses: a cold key on first access, a key lost to eviction, a key wiped by a failover.


Deep Dive 4: Write Policy and Invalidation - The Race That Poisons a Key Forever

Problem: The origin DB changes. The cache does not know. How you connect those two determines whether users see stale data for 300 ms, 300 seconds, or until someone notices and runs a manual DEL. There is a specific interleaving in the most common pattern that produces a permanently wrong cache entry, and it is the detail worth getting right.

The three policies, with what each actually costs:

Policy Write path Cost
Cache-aside App writes DB, then deletes the cache key. Next reader refills. Simplest and most common. Has the race below. First read after every write is slow.
Write-through App writes cache, cache writes DB synchronously, then acknowledge. No stale window and the cache is always warm. Every write pays both latencies, and a cache outage blocks writes.
Write-behind App writes cache, acknowledge immediately, flush to DB asynchronously. Fastest writes and it absorbs bursts by coalescing repeated writes to one key. A node loss before flush loses committed-looking writes - unacceptable for anything you would call a transaction.

Cache-aside is the right default here because of the volatility premise: the cache is not authoritative, so the DB must be written directly, and write-through makes the cache a hard dependency for writes. Write-behind is for counters and view tallies where the loss is a number being slightly wrong.

Bad: Cache-aside with TTL-only invalidation. Write to Postgres, do nothing to the cache, let the 300-second TTL handle it. Guaranteed up to five minutes of serving the old value, with no bound you can tighten without shortening the TTL for every key and tanking your hit rate. For a price change or a permission revoke this is a defect, not a trade-off.

Good: Cache-aside with delete-after-write. Commit to Postgres, then DEL the key. The next reader misses and refills from the authoritative row. Stale window drops from 300 s to the few milliseconds between commit and delete. Two problems remain. If the process dies between commit and delete, the stale entry survives to its TTL - which is the argument for keeping a TTL even with explicit invalidation. And the race:

The race. Delete-after-write has an interleaving that produces an entry which is wrong until its TTL expires. It needs a reader and a writer overlapping, and it is not rare under load.

sequenceDiagram
    participant R as Reader
    participant C as Cache
    participant W as Writer
    participant DB as Origin DB

    R->>C: GET product:98765
    C-->>R: MISS
    R->>DB: SELECT price
    DB-->>R: price = 100 old value
    Note over R: Reader now holds 100 but has not written it yet
    W->>DB: UPDATE price = 120
    DB-->>W: committed
    W->>C: DEL product:98765
    C-->>W: key absent nothing to delete
    R->>C: SET product:98765 = 100
    Note over C: Cache now holds 100 while the DB holds 120
    Note over C: Wrong until the TTL fires

Walk the ordering: the reader reads the old value from the DB, the writer then updates the DB and deletes a cache key that is not there yet, and finally the slow reader writes its now-obsolete value in. The delete happened before the set it was supposed to cancel. Nothing is retried, nothing errors, every component behaved correctly - and the cache is wrong until the TTL saves it.

The window is the gap between the reader’s DB read and its cache set, typically 1-20 ms. At 50K reads/sec against a key that gets written occasionally, you will hit it. The reason it is the most interview-relevant detail here is that it is invisible in code review: the delete-after-write pattern looks obviously correct, and the bug only exists in the interleaving.

Great: Three answers, in increasing strength. Pick based on how much a stale read costs.

1. Delete-after-write plus a short TTL as a backstop. Keep delete-after-write for the common case and accept the race, bounding its damage with a TTL you can defend - 30-60 s on mutable entities instead of 300. You still serve a wrong value, but for a bounded and short time, and you have not added machinery. This is what most systems ship, and it is the right call when a stale product description for 30 seconds costs nothing.

2. Versioned keys, which make the race structurally impossible. Carry the row’s version (a monotonic counter or updated_at) in the key:

product:98765:v7     the current entry
product:98765:v8     after the update

A writer that bumps the row to v8 writes a new key; the stale reader’s SET product:98765:v7 = 100 lands on a key nobody will ask for again. Readers need the current version, which they get from the row they are already fetching, or from a tiny separate pointer key. The entry is immutable once written, which is the property that kills the race - you are never overwriting, so ordering stops mattering. The cost is garbage: v7 occupies memory until its TTL expires. With a short TTL and LRU underneath, that is acceptable and the correctness is absolute.

3. Delete twice, with a delay. Delete the key after the commit, then schedule a second delete 500 ms later (a delay queue, or a timer). The second delete lands after any in-flight reader’s set. It is a shorter path to correctness than versioning, and it is probabilistic - a reader stalled for 600 ms still wins. Use it when versioning is too invasive to retrofit.

What about the write itself failing after the delete? Do the delete after the commit, never before. Delete-then-write means a reader can refill from the pre-write DB state in the gap, and you are back to a stale entry with no race needed. Commit first, invalidate second, always.

For a change feed instead of app-level deletes: tail the DB’s replication log with CDC (Debezium, DynamoDB Streams) and emit invalidations from the committed log. Now invalidation cannot be forgotten by a code path that writes the DB directly, which is the other common source of permanent staleness: a batch job or an admin tool that updates a row and knows nothing about your cache. CDC-driven invalidation is strictly more reliable than app-driven and costs you a pipeline to operate plus tens of milliseconds of lag.

In simple terms: The normal pattern is to update the database then delete the cached copy so the next reader refetches it. The subtle bug: a reader that already fetched the old value from the database, but has not yet saved it to the cache, will save that old value after your delete - and the cache is now wrong until its timer runs out. Fix it by never overwriting an entry (put the row’s version number in the key, so stale writers write a key nobody reads), or settle for a short timer that limits how long the wrong value survives.

Why versioned keys beat double-delete: Double-delete narrows the window; versioning removes it. Versioning makes cache entries immutable, and immutable entries have no write-ordering problem by construction. Double-delete is still a race, just one you are likely to win. Choose versioning for prices, balances, and permissions. Choose double-delete or a short TTL for descriptions, images, and feeds.


13. Design Self-Audit

Question Answer
Cold start after a full cluster restart? Hit rate is 0% and all 1M ops/sec hit an origin sized for 50K. The cache must never restart fully all at once - restart node by node, waiting for each to warm. If you genuinely need to cold-start the tier, you must shed load first: enable an admission limiter at the app so only a fraction of requests may fill, let the hit rate climb, then ramp. Alternative: a warming job that replays the top-N keys from yesterday’s hotkey report before taking traffic.
Can you lose data? Yes, by design, in four ways: eviction under memory pressure, TTL expiry, a failover dropping the async replication backlog, and a node crash with no replica ready. The origin DB is the only source of truth and must tolerate every cached key vanishing at once - which means the origin has to be sized for some multiple of steady-state miss traffic, not for 5%. If your system cannot function with a cold cache, you have built a database with no durability.
How does resharding avoid dropping requests? Slots migrate one at a time. While slot 9642 is moving, the source node still owns it and serves keys it still holds; for a key already migrated it replies ASK <target>, a one-shot redirect the client follows without updating its topology. Only when the slot is fully moved does ownership flip and MOVED start being returned, which is what tells clients to refresh the map. No request is dropped and no key is unreachable - the cost is an extra hop for in-flight keys during migration.
What if the origin DB is down? The cache keeps serving hits, so a read-only slice of traffic survives - this is the most underrated availability property of a cache tier. Misses fail. Three things make this materially better: stale-while-revalidate with a long stale_until so expiring keys keep serving, a circuit breaker so misses fail fast instead of each burning a 10-second DB timeout, and negative caching of the failure itself for 1-5 s so a dead origin is not hammered by every miss. Do not extend TTLs dynamically on origin failure - you will serve indefinitely stale data and forget you did.
Single points of failure? No proxy tier, so no shared data-plane component. Each shard is a primary plus a replica, and the gossip layer needs a majority of primaries for a failover vote. The real SPOF is correlated failure: 16 primaries in one AZ means an AZ loss takes the tier. Spread primaries across AZs and keep each replica in a different AZ from its primary, accepting ~1 ms of cross-AZ replication lag and cross-AZ data transfer cost.
Hit rate regression you would not notice? A global 95% hides a lot. A deploy that changes a key prefix, or a serialisation change without a version suffix, silently drops one entity type to a 0% hit rate while the aggregate barely moves. Alert on hit rate per key prefix, and on evicted_keys rising - sustained evictions mean the working set has outgrown memory and your hit rate is about to fall whether or not anything is broken.
Cost at scale? 32 nodes Γ— 64 GB is ~2 TB of RAM. Cache-grade instances run roughly $0.10-0.20 per GB of RAM per month in reserved cloud pricing, so order $3-6K/month, plus cross-AZ transfer on replication. The comparison that justifies it: serving 909K reads/sec from Postgres instead would need read replicas in the dozens. The cache is the cheap part.

14. Core Flows

Flow: GET with a Hit, and a Miss with Single-Flight Fill

sequenceDiagram
    participant A as App Thread
    participant L as Near Cache
    participant C as Cache Primary
    participant DB as Origin DB

    A->>L: get product:98765
    alt Near cache hit
        L-->>A: value in ~100ns
    else Cache tier hit
        L->>C: GET product:98765
        C-->>L: value
        L-->>A: value in ~300us
    else Cache tier miss
        L->>C: GET product:98765
        C-->>L: MISS
        L->>C: SET lock NX EX 5
        C-->>L: lock acquired
        L->>DB: SELECT row
        DB-->>L: row
        L->>C: SET product:98765 EX 300
        L->>C: DEL lock
        L-->>A: value in ~3ms
    end

Walkthrough:

  1. App asks its in-process Caffeine cache. Hit returns in ~100 ns with no syscall - for a hot key this is where most reads end.
  2. Near-cache miss goes to the cache tier. The client hashes the key to slot 9642, finds the owning primary in its cached slot map, and sends GET over a pooled connection.
  3. A hit returns in ~300 Β΅s, dominated by the network round trip. The client populates its near cache with a short TTL.
  4. On a miss the client attempts SET lock:product:98765 <owner> NX EX 5. Exactly one caller across the fleet wins.
  5. The winner queries Postgres, writes the value with a jittered TTL, and releases the lock with a compare-and-delete on the owner token.
  6. Losers poll the key every 20 ms up to a bound (say 3 attempts), then give up and query the origin directly rather than waiting forever.

Non-obvious failure: the lock holder stalls - a GC pause, a slow query, a dropped packet - and its 5-second lock expires while it is still working. A second caller acquires the lock and fetches too, so you get two origin queries instead of one. Harmless. The damaging version is when the first caller then finishes and runs DEL lock:product:98765, deleting the second caller’s lock and letting a third in. That is unbounded fan-out from a single stall. The fix is the compare-and-delete in step 5: the lock value holds the owner’s token and release only succeeds if the token still matches, which has to be a Lua script because check-then-delete is not atomic over two round trips.

The other failure worth naming: the losers’ fallback in step 6. Without the bound they wait on a holder that may never return. With the bound they query the origin directly, which means a pathological case degrades to the stampede you were preventing - bounded, but present. That is the right trade: a slow response is recoverable, a wedged request thread is not.

Flow: Primary Fails and a Replica Is Promoted

sequenceDiagram
    participant A as App
    participant P as Primary 10
    participant R as Replica 10
    participant G as Peer Primaries

    A->>P: GET key in slot 9642
    Note over P: Node crashes
    A->>P: GET times out after 200ms
    A->>A: Treat as miss and read origin
    G->>P: PING no response
    G->>G: Mark PFAIL locally
    G->>G: Majority agrees mark FAIL
    R->>G: Request promotion vote
    G-->>R: Vote granted
    R->>R: Promote to primary claim slots
    A->>P: GET key
    P-->>A: connection refused
    A->>A: Refresh CLUSTER SLOTS
    A->>R: GET key in slot 9642
    R-->>A: value or MISS

Walkthrough:

  1. The primary dies mid-request. The client’s 200 ms timeout fires - short on purpose, because a cache operation that takes longer than the origin query has stopped being a cache.
  2. That timeout is treated as a miss and the origin is read instead, so the user’s request still completes, just slower. A circuit breaker opens after a handful of consecutive failures so subsequent requests skip the dead node entirely instead of each paying 200 ms.
  3. Peer primaries stop getting PING replies and mark the node PFAIL - a local suspicion that carries no authority on its own.
  4. Suspicions gossip across the cluster bus. Once a majority of primaries report PFAIL for the same node it becomes FAIL cluster-wide. The majority requirement is what prevents one node with a broken NIC from declaring its healthy peers dead.
  5. The replica requests a promotion vote, wins a majority, promotes itself, and claims slots 0-1023 under a new config epoch. The epoch is how a stale old primary rejoining is recognised as out of date and demoted rather than accepted as a second owner.
  6. Clients refresh CLUSTER SLOTS on the first error from the old address and route to the new primary. Keys written in the last few hundred milliseconds before the crash may be missing - they were still in the replication backlog. Those reads miss, refill from the origin, and nobody notices.

Non-obvious failure: the primary is not dead, it is unreachable from the majority - a network partition with the primary on the minority side. The replica is promoted and starts taking writes while the old primary, still alive and still reachable from some app servers in its partition, keeps serving reads and accepting writes. For a few seconds you have two primaries for the same slots and clients in each partition see a different version of those keys. In a durable store that is a split-brain incident requiring reconciliation. In a cache it costs staleness for the length of the partition, and the config-epoch rule cleans it up: when the old primary reconnects it sees a higher epoch, demotes itself to replica, and discards its divergent data. The writes it accepted during the partition are simply gone. That is the volatility premise collecting its bill, and the reason this design is acceptable only because the origin DB holds the truth.


15. Final Architecture

flowchart LR
    APP["App Server<br/>client library plus near cache"]:::client
    TOPO["Slot Map<br/>16384 slots"]:::service
    P1["Primary 1<br/>slots 0 to 1023"]:::data
    P2["Primary 2<br/>slots 1024 to 2047"]:::data
    PN["Primary 16<br/>slots 15360 to 16383"]:::data
    R1["Replica 1"]:::data
    R2["Replica 2"]:::data
    RN["Replica 16"]:::data
    GOSSIP["Cluster Bus<br/>gossip and failover"]:::async
    HOT["Hot Key Detector<br/>top K sketch"]:::async
    METRICS["Prometheus<br/>hit rate evictions p99"]:::async
    DB[("Origin DB<br/>Postgres")]:::data

    APP -->|"1. Hash key to slot"| TOPO
    TOPO -->|"2. Route to owner"| P1
    TOPO -->|"3. Route to owner"| P2
    TOPO -->|"4. Route to owner"| PN
    P1 -->|"5. Async replicate"| R1
    P2 -->|"6. Async replicate"| R2
    PN -->|"7. Async replicate"| RN
    GOSSIP -->|"8. Detect failure and promote"| R1
    APP -->|"9. Single flight fill on miss"| DB
    APP -->|"10. Report key frequencies"| HOT
    HOT -->|"11. Replicate hot keys"| P1
    P1 -->|"12. Scrape stats"| METRICS

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
    classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0

How it works end-to-end (read path):

  1. App checks its near cache first β€” Caffeine, bounded to ~10K entries with a 1-5 second TTL. Absorbs hot keys entirely, ~100 ns per hit.
  2. Client hashes the key to a slot β€” CRC16(key) mod 16384, then the locally cached slot map resolves the owning primary. No proxy hop.
  3. Primary answers from RAM β€” hash table lookup, LRU or LFU metadata touched, value on the wire. ~300 Β΅s round trip inside one AZ.
  4. On a miss, single-flight fills β€” one caller across the fleet wins a SET NX lock, queries Postgres, writes the value with a jittered TTL, and releases the lock with a token-checked Lua delete.
  5. Early recomputation prevents most misses β€” readers probabilistically refresh a key shortly before its expiry, so hot keys are replaced while still valid and the herd never forms.

How it works end-to-end (write path and operations):

  1. Writes go to Postgres first, then invalidate β€” commit, then DEL the key. Mutable entities use a version suffix so a slow reader cannot write a stale value over a fresh one.
  2. Replication shadows every primary β€” async, so writes never wait. A failover loses the backlog, which costs a burst of misses and nothing else.
  3. Gossip handles membership and failover β€” a majority of primaries must agree a node is dead before a replica is promoted, and config epochs resolve a returning old primary.
  4. The hot key detector closes the loop β€” client-side frequency sketches report top-K keys; flagged keys get replicated under suffixed names across slots so one celebrity key cannot saturate one node.
  5. Metrics drive the operational decisions β€” hit rate per key prefix, evicted_keys, p99 per node, and memory fragmentation. Sustained evictions mean the working set has outgrown the tier.

Key Technologies

Term What it is
Consistent Hashing Maps keys and nodes onto a ring so adding a node moves ~1 of N keys instead of nearly all of them. The reason you can scale a cache without a cold-start outage.
Virtual Nodes 100-200 ring positions per physical node. Averages away the uneven arc lengths you get when each node holds a single position.
Hash Slot Redis Cluster’s 16384 fixed buckets. Keys hash to a slot permanently; only the slot-to-node assignment moves during resharding.
Approximated LRU Sample a few random keys and evict the stalest. Redis’s default maxmemory-samples is 5. Nearly true-LRU hit rate without 16 bytes of pointers per entry.
W-TinyLFU Caffeine’s admission policy. A frequency sketch decides whether a new entry deserves to displace an existing one, which makes the cache scan-resistant.
Single-Flight One caller fetches, the rest wait on its result. In-process it is a future map; across a fleet it is a SET NX lock with an owner token.
Stale-While-Revalidate Serve the expired value and refresh in the background. Keeps origin latency off the user path and keeps the cache useful when the origin is down.
MOVED and ASK Redis Cluster redirects. MOVED means ownership changed permanently, refresh your map. ASK is a one-shot redirect during an in-progress slot migration.
Hash Tag Wrapping part of a key in {...} so only that substring is hashed, forcing related keys into one slot for multi-key reads. Also a hand-made hot spot.

What’s Expected at Each Level

Mid-level

Get the basics right and in the right order. Explain why an in-memory cache exists at all with real numbers (RAM ~100 ns, SSD ~100 Β΅s, a DB query ~1 ms). Design the single node as a hash map plus a doubly linked list for LRU. Know that hash(key) % N breaks when N changes and that consistent hashing is the fix. Propose cache-aside with a TTL, and say out loud that cache data is disposable so the origin must handle a miss. Knowing LRU, FIFO, and TTL expiry as distinct ideas is table stakes.

Senior

Quantify the % N failure rather than asserting it - 90% of keys remapped on a 10-to-11 node change, and what that does to origin QPS in sequence. Explain virtual nodes with the variance argument for why 100-200 is the right range. Design replication as explicitly async and state what a failover loses. Handle hot keys with a near cache plus key replication, and name the staleness each one buys. Know that true LRU is not what production systems run and why sampling wins. Articulate the cache-aside stale-write race precisely - the interleaving, not just β€œthere can be staleness” - and propose delete-after-write with a TTL backstop.

Staff+

Drive the conversation to the uncomfortable parts. The cold-start problem: you cannot restart this tier all at once, and if you must, you need load shedding on the fill path, which means the app needs an admission limiter it probably does not have. Resharding without dropped requests, using ASK for in-flight keys versus MOVED for completed migration. Early recomputation with the mechanism, and why it beats locking (prevention versus damage control, and no waiter pays the fetch latency). Versioned keys as a structural fix for the stale-write race rather than a narrower window. CDC-driven invalidation to cover write paths that bypass the app. AZ topology and the correlated-failure argument against packing primaries into one zone. And the cost framing: the cache tier is cheap relative to the read replicas you would otherwise buy, so the interesting constraint is hit rate, not dollars.


🎯 Key Takeaways



Understand the building blocks used in this design:

Discussion

Newest first
You

Free system design + DSA prep. If it helped you crack an interview, consider supporting.

SensAI SensAI
Beta
Listening...
Tap mic to stop voice mode

Shape what we build next

Every piece of feedback is read by the team and directly influences our roadmap.

What type of feedback?

Install SystemCraft

Add to your home screen for instant access, offline reading, and a distraction-free experience.

Offline reading Faster loads No browser tabs App-like feel

Unlock AI Features

One click to activate - no payment, no credit card. Just sign in and you're in.

AI code review and hints
SensAI chat assistant
AI mock interviews
Whiteboard analysis
100% free during early access