Why this problem matters
When you type "ipl score" into a search engine, the engine has to look through an index of billions of documents, score candidates and build a results page. That work is expensive: it touches many index shards and takes tens to hundreds of milliseconds of machine time. Yet millions of people type the same few thousand queries every day. Running the full search again for each of them wastes enormous amounts of computing power.
A query cache stores the results page for a query the first time it is computed and serves later identical queries from memory. This lesson designs that cache: a key-value cache (a store that maps a key, here the query, to a value, here the results) sitting in front of a search engine. Then it generalizes the design into a distributed cache, which is the same machinery behind Redis clusters, Memcached fleets and the caches in almost every large system.
Interviewers use this problem to probe:
- whether you can implement LRU (least recently used) eviction and explain why it fits,
- memory math: how many entries, how big each one is, how many machines,
- how to split the cache across machines (sharding) and why plain modulo hashing fails,
- consistent hashing, replication and what happens when a node dies,
- invalidation: what to do when results change,
- hot keys and thundering herds.
The general caching patterns are in Caching; this lesson applies them end to end. The LRU data structure also appears in the low-level design lesson Design an LRU cache and rate limiter.
Step 1: Problem and scope
Build a cache that stores search results for popular queries so that repeated queries are answered from memory in a few milliseconds, the search backend does far less work, and results are never unacceptably stale.
In scope: the cache key, what to store, LRU eviction, sizing, sharding across many nodes, expiry and invalidation, replication and failure handling, hot keys, and the generalization to any distributed cache.
Out of scope: how the search engine ranks documents (it is a black box that takes a query and returns results), autocomplete (see Design typeahead autocomplete), personalization in depth.
Step 2: Clarifying questions
| Question | Assumed answer | Why it matters |
|---|---|---|
| Queries per month? | 10 billion | Request rate |
| Are results personalized? | Mostly not; location and language vary | Shapes the cache key |
| How fresh must results be? | Minutes for most queries, seconds for news-like queries | TTL and invalidation |
| What is cached: ids or full page? | Rendered result list: titles, URLs, snippets | Entry size |
| How many results per entry? | First page of 10 | Entry size |
| Is a miss acceptable? | Yes, the backend always answers | Cache is an optimization, not the source of truth |
| How popular are queries? | Heavily skewed; a small share of queries make up most traffic | Achievable hit ratio |
| Latency target for a hit? | Under 10 ms inside the data center | In-memory store, no disk |
The most important answer: the cache is not the source of truth. If an entry is lost, the backend recomputes it. That makes many failure cases simple: losing data costs speed, not correctness.
Step 3: Requirements
Functional requirements
get(query)returns cached results or reports a miss.set(query, results, ttl)stores results after the backend computes them.- Evict the least valuable entries when memory is full.
- Expire entries after their TTL, and support explicit invalidation.
Non-functional requirements
- Low latency: a hit in a few milliseconds, well below the backend's latency.
- High hit ratio: the share of queries served from cache. Every percentage point is backend capacity saved.
- High availability: a cache node failure should reduce the hit ratio briefly, not cause errors or overload the backend.
- Bounded staleness: results are allowed to be old by at most the TTL, and important changes (a page removed for legal reasons) must be purged quickly.
- Horizontal scalability: add memory by adding machines, without a full reshuffle.
Step 4: Back-of-the-envelope estimates
Traffic
- 10 billion queries per month ÷ 2,592,000 seconds ≈ 3,858 queries per second on average.
- Peak at 3 times average: about 11,600 queries per second.
Entry size
| Part | Approx. size |
|---|---|
| Key: normalized query plus language and region | 50 bytes |
| Value: 10 results × (title, URL, snippet) at about 250 bytes each | 2,500 bytes |
| Metadata (TTL, timestamps) and memory overhead of the store | about 450 bytes |
| Total | about 3 KB |
How many entries, and what hit ratio?
Query popularity is very skewed. A common model is a Zipf distribution: the query at popularity rank r is searched in proportion to 1/r. The top query is twice as frequent as the second, three times the third, and so on. Real query logs are roughly like this, though the exact shape varies, so treat the numbers below as an illustration of the method.
Under a Zipf model with exponent 1 over N distinct queries, caching the top C queries gives a hit ratio of about H(C) / H(N), where H(n) is the harmonic number 1 + 1/2 + ... + 1/n ≈ ln(n) + 0.5772.
With 1 billion distinct queries (N = 10 to the power 9, H(N) ≈ 21.30):
| Cached entries (C) | H(C) | Hit ratio H(C)/H(N) | Memory at 3 KB |
|---|---|---|---|
| 1 million | 14.39 | about 67.6% | 3 GB |
| 10 million | 16.70 | about 78.4% | 30 GB |
| 100 million | 18.99 | about 89.2% | 300 GB |
| 500 million | 20.61 | about 96.7% | 1.5 TB |
Step through one row to check it: for C = 100 million, H(C) ≈ ln(100,000,000) + 0.5772 = 18.42 + 0.58 = 18.99. Divide by 21.30 to get 0.892, an 89.2% hit ratio.
Two lessons from the table:
- Diminishing returns. Going from 1 million to 10 million entries (10 times the memory) gains about 11 points. Going from 100 million to 500 million (5 times the memory) gains only about 7.5 points.
- The long tail never fits. Half the queries are rare, and many queries are seen only once ever. No cache catches those.
We will size for 100 million entries, about 300 GB, targeting roughly a 90% hit ratio. (LRU keeps approximately the most popular entries, not exactly, so the real hit ratio will be a little lower than this ideal.)
Machines
- Suppose each cache node has 64 GB of RAM, of which about 48 GB is usable for data after the operating system, fragmentation and headroom.
- 300 GB ÷ 48 GB = 6.25, so 7 primary nodes. With one replica each, 14 nodes. Round up to leave room for growth.
Backend savings
- Without a cache, the backend serves all 11,600 peak queries per second.
- With an 89.2% hit ratio, it serves 10.8%: about 1,250 per second. The cache cut backend load by roughly 9 times. That is the business case for the whole system.
What the numbers told us
About 11,600 queries per second at peak, 3 KB per entry, and 100 million entries for roughly a 90% hit ratio: about 300 GB of memory over 7 shards, doubled for replicas. The backend load drops about 9 times.
Step 5: API design
The cache is internal, used by the search front-end service.
GET /cache/{key} 200 {value, expiresAt} | 404 miss
PUT /cache/{key} {value, ttlSeconds} 204
DELETE /cache/{key} 204
(In practice this is Redis or Memcached's own protocol:
GET key / SET key value EX ttl / DEL key)
The cache key
The key must be identical for queries that should share results and different for queries that should not. Build it from a normalized query plus every input that changes the results:
key = hash( normalize(query) + "|" + language + "|" + region
+ "|" + safeSearch + "|" + page )
normalize(" IPL Score ") -> "ipl score"
key("ipl score", en, IN, strict, 1) -> "q:7f3a9c..."
Normalization lowercases, trims and collapses spaces, and possibly removes punctuation. Be careful: lowercasing is safe for most text, but some normalization (removing accents, reordering words) can change meaning, so normalize only what the search engine itself treats as equivalent.
Common mistake
Forgetting an input in the key. If region is not part of the key, a user in India searching "weather" might receive results cached for a user in London. Anything that changes the results must be in the key, and personal data must never be in a shared cache entry.
Step 6: Data model and storage choice
Each cache node stores entries in memory:
entry
key string (about 50 bytes, or a 16-byte hash)
value bytes (serialized, often compressed results)
expires_at int64 (Unix time)
lru links prev / next pointers (LRU list)
Store choice.
- Redis: in-memory, single-threaded per instance for commands (newer versions use extra threads for network I/O), rich data types, optional persistence and replication, cluster mode with 16,384 hash slots. Eviction policies include
allkeys-lruandallkeys-lfu, which use approximate sampling rather than a perfect LRU list. - Memcached: simpler, multi-threaded, pure key-value, no persistence or built-in replication. Good for a plain cache.
- In-process cache (a map inside each search server): fastest, but small and duplicated on every server. Useful as a tiny first layer for the very hottest queries.
Either Redis or Memcached is a sound answer for the shared layer. Persistence is unnecessary because the backend can rebuild everything.
Step 7: High-level design
user query
|
+-----v------+
| Load |
| balancer |
+-----+------+
|
+-----v------------------+
| Search front-end |
| (stateless) |
| - normalize, build key|
| - small local cache |
| - cache client with |
| hash ring |
+--+------------------+--+
| 1. get(key) | 3. on miss: run search
v v
+--+---------------+ +---------------------+
| Distributed | | Search backend |
| cache cluster | | (index shards, |
| | | rankers) |
| shard 1 -> repl | +---------------------+
| shard 2 -> repl |
| ... |<-- 4. set(key, results, ttl)
| shard 7 -> repl |
+------------------+
^
| invalidation events
+--------+---------+
| Index updater / |
| purge service |
+------------------+
The front-end is stateless. Its cache client holds the hash ring (the map of which node owns which keys), so it talks to the right shard directly with no proxy hop. Some deployments put a proxy (such as Twemproxy or Envoy) in between instead, which centralizes routing at the cost of an extra hop.
Step 8: Request flows
Read path (cache hit)
- The front-end receives "IPL Score" from a user in India.
- It normalizes to "ipl score" and builds the key with language, region, safe-search setting and page number.
- It checks its small in-process cache. Suppose a miss.
- It hashes the key onto the ring, finds the owning node, and sends
GET key. - The node finds the entry, checks it has not expired, moves it to the front of its LRU list, and returns it.
- The front-end returns the results. Total cache time: about a millisecond or two.
Read path (cache miss) and write
- Steps 1 to 4 as above; the node replies "miss".
- The front-end sends the query to the search backend, which takes, say, 100 to 300 ms.
- The front-end returns results to the user and sends
SET key value EX ttlto the cache node. This write can be asynchronous so the user does not wait for it. - If memory is full, the node evicts its least recently used entry to make room.
This is the cache-aside pattern (also called lazy loading): the application checks the cache, falls back to the source, then fills the cache. Only queries that are actually asked get cached, which suits a skewed workload perfectly.
Step 9: Deep dives
Deep dive 1: LRU eviction and how to implement it
Memory is finite, so when the cache is full, something must go. LRU evicts the entry that was used least recently. The bet is that a query asked a minute ago is more likely to be asked again soon than one last asked yesterday. For search, where popularity is skewed and shifts over time (news, sports, festivals), this works well.
An LRU cache needs both operations in O(1) time:
- Look up by key: a hash map from key to entry.
- Know the least recent entry and move any entry to "most recent": a doubly linked list ordered by recency. The head is the most recently used; the tail is the least. Because each node has both
prevandnextpointers, you can unlink any node in O(1) once the hash map has found it.
hash map: "ipl score" -> node B, "gold price" -> node A, ...
head <-> [A: gold price] <-> [B: ipl score] <-> [C: weather] <-> tail
most recent least recent
get("ipl score"): unlink B, insert at head
head <-> [B] <-> [A] <-> [C] <-> tail
put("train status") when full: remove C (tail side), insert new at head
class Node:
__slots__ = ("key", "value", "prev", "next")
def __init__(self, key=None, value=None):
self.key, self.value = key, value
self.prev = self.next = None
class LRUCache:
def __init__(self, capacity: int):
self.capacity = capacity
self.map = {} # key -> Node
self.head, self.tail = Node(), Node() # sentinels
self.head.next, self.tail.prev = self.tail, self.head
def _unlink(self, node):
node.prev.next, node.next.prev = node.next, node.prev
def _push_front(self, node):
node.prev, node.next = self.head, self.head.next
self.head.next.prev = node
self.head.next = node
def get(self, key):
node = self.map.get(key)
if node is None:
return None # miss
self._unlink(node)
self._push_front(node) # now most recently used
return node.value
def put(self, key, value):
node = self.map.get(key)
if node:
node.value = value
self._unlink(node)
else:
if len(self.map) == self.capacity:
lru = self.tail.prev # least recently used
self._unlink(lru)
del self.map[lru.key]
node = Node(key, value)
self.map[key] = node
self._push_front(node)
cache = LRUCache(capacity=2)
cache.put("ipl score", ["r1", "r2"])
cache.put("weather delhi", ["r3"])
cache.get("ipl score") # touch: now most recent
cache.put("gold price", ["r4"]) # evicts "weather delhi"
print(cache.get("weather delhi"), cache.get("ipl score"))
# None ['r1', 'r2']
The two sentinel nodes (dummy head and tail) mean you never have to special-case an empty list or the first and last elements. In Python, collections.OrderedDict with move_to_end gives the same behavior in a few lines; in Java, LinkedHashMap with access order and removeEldestEntry does too. Interviewers often want the hand-built version.
LRU at scale. A perfect LRU list needs every read to update shared pointers, which requires a lock in a multi-threaded server and costs two pointers per entry. Large caches therefore approximate:
- Sharded LRU: split the cache into many independent segments, each with its own lock and list, so threads rarely contend.
- Sampled LRU (Redis's approach): store a last-access timestamp per key; when memory is full, sample a handful of keys and evict the oldest among them. Nearly as good, far cheaper.
- CLOCK: a circular buffer with a "recently used" bit per entry; a hand sweeps around clearing bits and evicts the first entry whose bit is already clear.
LRU versus LFU. LFU (least frequently used) evicts entries with the lowest access count. It protects consistently popular queries from being pushed out by a burst of one-time queries (a "scan"), but adapts slowly when popularity shifts, unless counts decay over time. Modern designs such as TinyLFU and W-TinyLFU (used by the Java Caffeine library) combine a frequency sketch with LRU to get the strengths of both. A good answer: "LRU is the default; if one-hit-wonder queries pollute the cache, add an admission filter that only caches a query once it has been seen twice, or use an LFU-style policy."
Interview tip
When asked to implement LRU, say the data structures first: "a hash map for O(1) lookup and a doubly linked list for O(1) reordering and eviction, with sentinel head and tail nodes." Then write get and put. Mention thread safety and approximate LRU as follow-ups before you are asked.
Deep dive 2: Sharding the cache and consistent hashing
300 GB does not fit on one node, so keys must be divided among 7 nodes. Every front-end must agree on which node owns which key without asking anyone.
Naive approach: modulo hashing. node = hash(key) mod N. It spreads keys evenly, but changing N breaks it. If you go from 4 nodes to 5, a key stays on the same node only when hash mod 4 equals hash mod 5, which happens for 4 out of every 20 hash values. So 80% of keys move. For a cache, a moved key is a miss: the hit ratio collapses, and the backend suddenly faces most of the traffic. Adding capacity causes an outage.
Consistent hashing. Picture a ring of hash values from 0 to 2 to the power 64 minus 1, wrapping around.
- Hash each node's name to a point on the ring.
- Hash each key to a point on the ring.
- A key belongs to the first node found by walking clockwise from the key's point.
Ring positions, clockwise from 0 (wrapping at 2^64):
0 ---- A(10) ---- k1(25) ---- B(40) ---- k2(55)
---- C(70) ---- k3(85) ---- D(95) ---- back to 0
k1 at 25 -> next node clockwise is B(40) -> owner B
k2 at 55 -> next node clockwise is C(70) -> owner C
k3 at 85 -> next node clockwise is D(95) -> owner D
Add node E at 50: only keys between 40 and 50 move (B -> E).
k1, k2 and k3 all keep their owners.
When a node is added, it takes over only the keys between its point and the previous node's point; every other key stays put. With N nodes, adding one moves about 1/(N+1) of the keys. When a node is removed, only its keys move, to its clockwise neighbor.
Virtual nodes. With only 4 or 5 points on the ring, the arcs between them are very uneven, so one node might own twice as many keys as another. The fix: place each physical node at many points (say 100 to 200 "virtual nodes"), using hash("cache-a#0"), hash("cache-a#1"), and so on. Arcs average out, load evens out, and when a node leaves, its keys spread across all remaining nodes rather than landing on one neighbor. Virtual nodes also allow weighting: a bigger machine gets more points.
import bisect
import hashlib
def h(key: str) -> int:
return int.from_bytes(hashlib.md5(key.encode()).digest()[:8], "big")
class HashRing:
def __init__(self, nodes, vnodes: int = 100):
self.vnodes = vnodes
self.ring = [] # sorted list of (hash, node)
for node in nodes:
self.add(node)
def add(self, node: str) -> None:
for i in range(self.vnodes):
bisect.insort(self.ring, (h(f"{node}#{i}"), node))
def remove(self, node: str) -> None:
self.ring = [(p, n) for p, n in self.ring if n != node]
def node_for(self, key: str) -> str:
i = bisect.bisect(self.ring, (h(key),)) # first point clockwise
return self.ring[i % len(self.ring)][1]
keys = [f"query-{i}" for i in range(100_000)]
ring = HashRing(["cache-a", "cache-b", "cache-c", "cache-d"])
before = {k: ring.node_for(k) for k in keys}
ring.add("cache-e")
moved = sum(before[k] != ring.node_for(k) for k in keys)
print(f"consistent hashing moved {moved / len(keys):.1%}") # about 20%
mod_moved = sum(h(k) % 4 != h(k) % 5 for k in keys)
print(f"modulo hashing moved {mod_moved / len(keys):.1%}") # about 80%
Running this moves about 19% of keys with consistent hashing (the ideal is 1/5 = 20%) and exactly 80% with modulo hashing. That contrast is the entire argument for consistent hashing.
Alternatives you should know. Redis Cluster uses a fixed set of 16,384 hash slots: slot = CRC16(key) mod 16384, and each node owns a range of slots. Adding a node means moving some slots, not rehashing everything, which gives the same benefit with explicit control. Rendezvous hashing (highest random weight) computes a score for each node per key and picks the highest; it also moves only about 1/N of keys and needs no ring. For more, see Databases at scale.
Deep dive 3: Expiry and invalidation
"There are only two hard things in computer science: cache invalidation and naming things" is an old joke because it is true. Search results change when pages are added, removed or re-ranked. How do cached results stay acceptable?
1. TTL (time-to-live) on every entry. The simplest and most important tool. Every entry expires after a set time, so staleness is bounded by the TTL. Choose TTLs per query type:
| Query type | Example | Suggested TTL |
|---|---|---|
| Evergreen, informational | "binary search algorithm" | hours |
| General | "best phones under 20000" | 15 to 60 minutes |
| News and events | "election results", "ipl score" | seconds to 1 minute |
A classifier, or simple signals such as "this query had a large spike in volume" or "results contain news sources", can pick the TTL.
2. Expiry implementation. Like Redis, combine lazy expiry (check expires_at on every read and treat expired entries as misses) with active expiry (a background task samples keys with TTLs and deletes expired ones, so dead entries do not occupy memory forever).
3. Explicit purges. Some changes cannot wait for a TTL: a page removed for legal reasons, malware discovered on a result, a privacy request. The purge service needs to find all cached queries whose results include that URL. Options:
- maintain a reverse index, URL to list of cache keys containing it, updated on every
set(accurate but costly to maintain), - filter at read time: keep a small, replicated "blocked URLs" set and strip those URLs from any cached result before returning it (cheap and immediate; this is the common practical answer),
- for wide changes, bump a version number in the key (
v42:...). All old entries become unreachable at once and age out through LRU.
4. Stale-while-revalidate. When a popular entry expires, serve the slightly old value to the current requester and refresh it in the background. Users see no latency spike, and the backend sees one refresh instead of thousands.
5. Index updates. When the search index publishes a new version (for example every few minutes), you could invalidate everything, but that would empty the cache and overload the backend. Better: let TTLs roll entries over gradually, and use version bumps only for categories of queries known to be affected.
Common mistake
Saying "we will update the cache whenever the data changes" for search. A single new web page can affect the results of millions of queries, and you cannot know which ones without running them. For search, bounded staleness through TTLs, plus targeted purges for urgent removals, is the realistic design.
Deep dive 4: Hot keys, thundering herds and replication
Hot keys. During a cricket final, "ipl score" might receive tens of thousands of requests per second. Consistent hashing sends all of them to the one node that owns that key. That node's network or CPU saturates while others sit idle. Fixes:
- Local (in-process) cache on each front-end with a very short TTL (1 to 5 seconds). The hottest keys are served without leaving the machine. With 100 front-ends and a 2-second TTL, the shared node sees at most about 50 requests per second for that key, no matter how many users ask.
- Key replication: store hot keys under several suffixes (
ipl score#1...ipl score#8), which land on different nodes, and have readers pick one at random. - Read replicas: let replicas of the owning shard serve reads.
- Detection: count requests per key with a small frequency sketch (for example a count-min sketch, a compact structure for approximate counts) and promote keys automatically when they cross a threshold.
Thundering herd (cache stampede). A hot entry expires, and in the next 100 ms, 2,000 requests all miss and all hit the backend with the same query. Fixes:
- Request coalescing: only the first miss for a key calls the backend; others wait for its result. Within one server this is a map from key to in-flight request; across servers, use a short lock in the cache (
SET lock:key NX EX 5). - Stale-while-revalidate, described above.
- Jittered TTLs: add a random few percent to each TTL so that entries filled at the same moment do not all expire at the same moment.
Replication and node failure. Each shard has a primary and a replica. Writes go to the primary and are copied asynchronously; if the primary fails, the replica is promoted (Redis Sentinel or Redis Cluster handles this). Because the cache is not the source of truth, asynchronous replication is fine: losing the last few writes costs a few misses, not wrong answers.
What happens to the backend when one of 7 shards is lost without a replica? About 1/7 of the keys miss until refilled. At peak that adds roughly 11,600 × 0.892 ÷ 7 ≈ 1,480 extra backend queries per second, more than doubling the 1,250 the backend normally sees. The backend must have that headroom, or the cache layer needs replicas, or the front-end must shed load. Calculating this explicitly is a strong interview move.
Step 10: Generalizing to a distributed cache
Strip away "search" and you have the classic question "design a distributed cache" (essentially, design Memcached or Redis Cluster). The same pieces apply:
| Concern | Design choice |
|---|---|
| Data per node | Hash map plus LRU (or approximate LRU/LFU) |
| Partitioning | Consistent hashing with virtual nodes, or hash slots |
| Routing | Smart client with the ring, or a proxy layer |
| Membership | Config service (ZooKeeper, etcd) or gossip protocol tells clients which nodes exist |
| Replication | Primary plus replicas per shard, asynchronous |
| Failover | Health checks; promote replica; update ring |
| Expiry | Lazy plus active expiry |
| Hot keys | Local cache, key replication, read replicas |
| Consistency with the database | Cache-aside with delete-on-write, write-through, or write-behind |
The last row matters for general caches, where the cache sits in front of a database that the application also writes to:
- Cache-aside with delete on write: the application updates the database, then deletes the cache key. The next read reloads it. Delete rather than update, because two concurrent updates can otherwise leave the older value in the cache.
- Write-through: writes go to the cache and the database together; reads are always warm, but every write pays both costs.
- Write-behind (write-back): writes go to the cache and are flushed to the database later; fast, but risks data loss if the cache node dies before flushing.
For a search query cache, there are no application writes to results, so cache-aside with TTLs is the clear fit. See Caching and Scalability for more on these patterns.
Step 11: Scaling and bottlenecks
- More memory: add nodes; consistent hashing moves about 1/(N+1) of keys, and those misses refill gradually. Add nodes during low traffic.
- More throughput: add replicas to serve reads, and use the local cache tier for the hottest keys.
- Network: a 3 KB value at 11,600 requests per second is about 35 MB/s across the cluster, which is fine; compressing values cuts it further.
- Front-end connections: use connection pooling and pipelining (batching several commands on one connection) to the cache nodes.
- Multiple regions: run an independent cache cluster per region. Popular queries differ by region anyway, and cross-region cache reads would add more latency than they save.
Step 12: Failure handling
| Failure | Effect | Mitigation |
|---|---|---|
| Cache node down | 1/N of keys miss, backend load rises | Replica promotion; backend headroom; load shedding |
| Whole cache cluster down | Every query hits the backend | Rate-limit and degrade (fewer results, skip expensive features); circuit breaker |
| Slow cache node | Front-end requests stall | Short timeouts (a few ms); treat timeout as a miss |
| Hot key | One node saturated | Local cache, key replication, read replicas |
| Mass expiry | Stampede on backend | Jittered TTLs, coalescing, stale-while-revalidate |
| Network partition between client groups | Clients disagree on ring membership | Central membership service; short-lived inconsistency is only extra misses |
| Bad results cached (bug) | Wrong results served until TTL | Version bump in key to drop everything at once |
A guiding rule: a cache failure must never be worse than having no cache. Timeouts are short, errors are treated as misses, and the backend is protected by rate limits. See Reliability and recovery.
Step 13: Trade-offs and alternatives
| Decision | Choice here | Alternative | When the alternative wins |
|---|---|---|---|
| Eviction | LRU (approximate) | LFU / W-TinyLFU | Many one-time queries pollute the cache |
| Partitioning | Consistent hashing, virtual nodes | Hash slots (Redis Cluster) | You use Redis Cluster's built-in routing |
| Routing | Smart client | Proxy | Many languages or clients; want central control |
| Freshness | TTL per query type + purges | Event-driven invalidation | Small key space with known dependencies |
| Replication | Async primary/replica | No replicas | Backend can absorb losing a shard |
| Value | Rendered results | Document ids only | Snippets change often or are personalized |
| Tiers | Local + shared | Shared only | Few hot keys, small front-end fleet |
What interviewers probe
"Why not just use modulo hashing?" Changing the node count moves most keys (80% when going from 4 to 5 nodes), which wipes out the hit ratio. Consistent hashing moves only about 1/(N+1).
"Why virtual nodes?" A few points on the ring produce uneven arcs and uneven load, and a failed node dumps all its keys onto one neighbor. Many virtual points per node average this out and spread the failed node's keys across everyone.
"How do you size the cache?" Entry size times the number of entries needed for a target hit ratio, divided by usable memory per node, then multiplied by the replication factor. Under our assumptions: 3 KB × 100 million = 300 GB, 7 nodes, 14 with replicas.
"What if results must be perfectly fresh?" Then caching full results is the wrong tool for those queries. Use very short TTLs, or cache only the expensive inner parts (index shard results) that change less, or skip the cache for that class of query.
"How would you cache personalized results?" Cache the non-personal part (candidate documents for the query) in the shared cache, and apply personalization on top per user. Never put user-specific data in a shared entry.
"LRU or LFU?" LRU adapts quickly to shifting popularity, which suits search trends. LFU keeps steady favorites but adapts slowly. Approximate hybrids such as W-TinyLFU are a strong modern answer.
Interview questions
Q1. What is a cache hit ratio and why does it matter so much here?
It is the fraction of requests answered from the cache. Every miss costs a full search on the backend, so the hit ratio directly sets backend capacity. Under our assumptions, an 89% hit ratio cut peak backend load from about 11,600 to about 1,250 queries per second.
Q2. How would you implement an LRU cache with O(1) operations?
Use a hash map from key to node and a doubly linked list ordered by recency. On get, find the node via the map, unlink it and move it to the head. On put, insert at the head and, if over capacity, remove the node before the tail sentinel and delete it from the map.
Q3. Why do large caches use approximate LRU?
Exact LRU must update shared list pointers on every read, which needs locking and costs memory per entry. Sampling a few keys and evicting the oldest (as Redis does), or using CLOCK or segmented LRU, gives nearly the same hit ratio with much less overhead.
Q4. How do you size the cache?
Estimate entry size (about 3 KB here), choose how many entries you need for a target hit ratio using the popularity distribution, and multiply. 100 million entries need about 300 GB; at 48 GB usable per node, that is 7 nodes, and 14 with one replica each.
Q5. Explain consistent hashing.
Nodes and keys are hashed onto the same ring, and each key belongs to the next node clockwise. Adding or removing a node only moves the keys in the affected arc, about 1/N of the total. Virtual nodes place each node at many points to even out load.
Q6. What happens with modulo hashing when you add a node?
Most keys change owner because hash mod N changes for most values; going from 4 to 5 nodes moves 80% of keys. In a cache, moved keys are misses, so the hit ratio collapses and the backend can be overwhelmed.
Q7. How do you keep cached search results fresh?
Give every entry a TTL suited to the query type: hours for evergreen queries, seconds for news. Use explicit purges or read-time filtering for urgent removals, version numbers in keys to drop whole groups, and stale-while-revalidate to refresh popular entries smoothly.
Q8. What is a cache stampede and how do you prevent it?
When a popular entry expires, many concurrent requests miss together and all hit the backend. Prevent it with request coalescing (one loader, others wait), stale-while-revalidate, and random jitter on TTLs so entries do not expire in sync.
Q9. How do you handle a hot key?
Add a short-TTL in-process cache on each front-end, replicate the key under several suffixes on different nodes, or serve reads from replicas. Detect hot keys automatically with approximate counting.
Q10. What should the cache key contain?
The normalized query plus every input that changes results: language, region, safe-search setting, page number and result-format version. Missing an input serves wrong results; including personal data in a shared key leaks privacy and destroys the hit ratio.
Q11. Does the cache need persistence or synchronous replication?
No. The cache is not the source of truth, so losing entries only costs misses. Asynchronous replication is enough to avoid a big drop in hit ratio when a node fails, and persistence is optional for faster warm-up.
Q12. What happens if a cache node fails?
Its keys miss until a replica is promoted or the ring is updated, and those queries go to the backend. Under our numbers, losing 1 of 7 shards adds about 1,480 backend queries per second at peak, so the backend needs headroom or load shedding.
Q13. Compare cache-aside, write-through and write-behind.
Cache-aside loads on a miss and deletes on writes, which is simple and lazy. Write-through writes the cache and database together, keeping the cache warm at the cost of slower writes. Write-behind writes the cache first and the database later, which is fast but can lose data if a node dies.
Q14. When is hash-slot partitioning preferred over a ring?
When you want explicit, operator-controlled placement, as in Redis Cluster's 16,384 slots: moving a node's capacity means moving specific slots. It gives the same "move only a little" property as consistent hashing with simpler reasoning about which node owns what.
Key takeaways
- A query cache trades memory for backend work; the hit ratio is the number that matters.
- Popularity is skewed, so a cache with a small fraction of all queries catches most traffic; under a Zipf model, 100 million of 1 billion queries gives about 89%.
- Sizing: entry size times entries, divided by usable memory per node, times replicas; here about 300 GB, 7 shards, 14 nodes.
- LRU uses a hash map plus a doubly linked list for O(1) get and put; big caches approximate it.
- Consistent hashing with virtual nodes moves only about 1/(N+1) of keys when a node is added; modulo hashing moves most of them.
- Freshness comes from per-query-type TTLs, targeted purges, key versions and stale-while-revalidate.
- Hot keys and stampedes are handled with local caches, key replication, request coalescing and jittered TTLs.
- The cache is never the source of truth, so failures should cost speed, never correctness.
Next lesson
Continue with Design a social graph.

