System design interview guide
Design a Distributed Cache: Sharding, Eviction, Hot Keys and Stampedes
Two courses by the author of this page:
770 lessons · 18 free to read
₹499 in India · $49 elsewhere, once
Get System DesignYou own this course
204 lessons · 10 free to read
₹999 in India · $49 elsewhere, once
Get AI EngineeringYou own this course
A cache is a fast copy of data kept in memory, so most reads never reach the slower database behind it. Take a busy app that reads 1,000,000 times a second at peak. If the cache answers 95% of those reads, the database sees 50,000 reads a second. If the cache answers only 80%, the database sees 200,000: four times more, from a drop of 15 points. One machine cannot hold or serve all of it, so the cache is spread over many machines. Most of this design is about what happens when those machines fill up, fail, get added, or when one key becomes far more popular than the rest.
A distributed cache is a set of machines that each hold part of the cached data in memory (RAM), in front of a database that stays the source of truth. Each item is a key and a value. The app usually follows cache-aside: read the cache first; on a miss, read the database and put the value in the cache; on a write, update the database and delete the cached key. Keys are spread over the cache machines with consistent hashing. A hash function turns a key into a number; each machine owns many points on a ring of hash values (virtual nodes), and a key belongs to the first point clockwise from its own hash. Adding a machine then moves only about 1/N of the keys. Each machine evicts (removes) old items when it is full, usually by LRU, least recently used, which a hash map plus a doubly linked list does in constant time. Every item also gets a TTL, a time to live, after which it expires. The hard parts are failures and load spikes. A dead machine turns all its keys into misses that hit the database at once. One very popular key can overload one machine. When a popular key expires, thousands of requests can miss together and flood the database; this is a cache stampede. Facebook's memcache paper fixes stampedes with leases: only one client per key is allowed to refill it, and the others wait briefly and retry. In their measurement, peak database reads for such keys fell from 17,000 to 1,300 a second.
Where it shows up
This is a common question in backend and infrastructure interviews, asked as 'design a distributed cache' or 'design Memcached'. It also sits inside bigger questions, such as a news feed, a URL shortener or a rate limiter, whenever the interviewer asks how reads stay fast. Several companies built this layer and wrote about it in detail. Facebook described its memcache setup in the NSDI 2013 paper 'Scaling Memcache at Facebook', and released mcrouter, its routing proxy. Twitter released twemproxy. Amazon's Dynamo paper made virtual nodes on a hash ring widely known, and the Redis Cluster specification documents another way to split keys.
Spec sheeta Distributed Cache
- 01Items and sizes
- 500 million items, about 1,100 bytes each
- 02Memory
- about 690 GB, so 14 nodes
- 03Reads per node
- about 71,400 a second
- 04Database reads behind the cache
- 50,000 a second at a 95% hit rate
worked through below, with the maths
Why this question is asked
The question looks simple, because everyone has used a cache. It gets hard fast. To pass, you have to explain how keys are spread over machines, and what happens when the machine list changes. You have to say what a full machine throws away, and why. You have to keep the cache and the database from disagreeing, and say how long they may disagree. The interviewer also wants to hear about failures: a dead node, a hot key, an expired key that thousands of users want at once. A strong candidate treats the cache as something that can hurt the database, not only protect it, and says how to stop that. Note that this page designs the cache tier as a whole. How one Redis server works inside (its single thread, its persistence files) is covered on the Design Redis page.
This page is the free part.
The course goes deeper on the Distributed Cache design
₹499 in India$49 everywhere elseonce, for the whole course
The System Design course covers the Distributed Cache design across a run of lessons, not one page. These 4 alone are about 99 minutes of step-by-step reading, every one with a quiz.
- Distributed Cache
How a distributed cache spreads keys across nodes with consistent hashing and virtual nodes. Measured tests show why a hot key still overloads a node.
- Consistent Hashing
The algorithm that makes adding or removing nodes cheap, an interview favorite and production essential.
- Cache Eviction Policies
LFU beat LRU by 6 to 15 points, FIFO measured identical to random, and the scan that supposedly destroys LRU cost it 11 points for four hundred requests.
- Cache-Aside Pattern
Cache-aside is said to degrade gracefully when the cache fails. Measured, unwrapped it answered 0.0 percent of requests, and wrapped it answered 13.0, which is exactly the database's capacity.
all part of
System Design Masterclass
770 lessons · about 283 hours · 18 free to read
₹499in India, by UPI
$49everywhere else, by PayPal
The concepts were explained in a clear and structured way, with practical examples that made even complex system design topics easier to understand. I especially liked the focus on real-world architecture, scalability, trade-offs. The content was well designed, engaging, and highly useful for anyone looking to strengthen their system design skills. Highly recommended for software engineers preparing for system design interviews or wanting to build a stronger foundation in designing scalable systems.
You own this course
Continue with Distributed Cachepay once,
yours for life
Live from the course bench
Before you design Distributed Cache, try the questions it rests on.
Real interview questions, live diagrams and measured lab results, pulled at random from the lessons. A new set every time you come back.
Requirements
Always clarify these in the first 5 minutes of the interview. Do not start drawing boxes until both lists are agreed.
Functional requirements
- get(key): return the value, or a miss if the cache does not have it
- set(key, value, ttl): store a value, with a time to live (TTL) in seconds after which it expires
- delete(key): remove a value, used when the database copy changes
- Spread keys over many cache nodes (machines), and let nodes be added or removed while the system runs
- Evict (remove) items when a node's memory is full, keeping the items most likely to be read again
- Optional: read many keys in one request (a multi-get), because one web page often needs hundreds of items
Non-functional requirements
- Low latency (time to answer): around a millisecond or less for a hit inside one data centre
- High throughput: around a million reads a second at peak, growing by adding nodes
- High hit rate: 95% or more of reads answered from the cache, so the database stays small
- Availability: one dead node must not take the site down or flood the database
- Bounded staleness: a cached value may be out of date for a short, known time, never forever
- The database stays the source of truth. Losing cached data is allowed; it can be read again
Back-of-envelope scale estimates
Show your math. Pulling numbers from thin air signals you have not thought about the load.
Items and sizes
500 million items, about 1,100 bytes each
Assume 500 million hot items, such as user profiles and post counts. Assume an average value of 1,000 bytes, plus about 100 bytes for the key and the cache's own bookkeeping (pointers, expiry time, flags). These are interview assumptions. Say them out loud before using them, and measure the real sizes later.
Memory
about 690 GB, so 14 nodes
500,000,000 x 1,100 bytes = 550 GB. Add 25% headroom: memory is handed out in fixed-size chunks, so part of each chunk is wasted, and the data grows. 550 x 1.25 = 687.5 GB. If each node has 64 GB of RAM and gives 50 GB to the cache, 687.5 / 50 = 13.75, so 14 nodes. With one replica (copy) of each node, that is 28 machines.
Reads per node
about 71,400 a second
1,000,000 reads a second / 14 nodes = 71,429 a second each, if keys are spread evenly. The memcached docs say one fast machine with a fast network can handle 200,000 or more requests a second, so plan for 100,000 per node and keep the rest as headroom. Uneven keys and hot keys eat that headroom first.
Database reads behind the cache
50,000 a second at a 95% hit rate
Misses go to the database: 1,000,000 x 5% = 50,000 reads a second. At an 80% hit rate it is 1,000,000 x 20% = 200,000. The hit rate is the number that sizes the database. Writes (assume 50,000 a second) go to the database either way.
Network
about 8.8 Gbit/s in total
1,000,000 reads x 1,100 bytes = 1.1 GB a second. 1.1 x 8 = 8.8 gigabits a second for the whole cluster, or 8.8 / 14 = 0.63 Gbit/s per node. A normal server network card handles this, so memory decides the node count here, not network.
One node dies, no replica
database reads go from 50,000 to about 118,000 a second
The dead node served 71,429 reads a second. Before, 5% of them missed: 3,571. Now all of them miss. 50,000 - 3,571 + 71,429 = 117,858 database reads a second, 2.4 times normal, until the keys are cached again somewhere. This one number is why the failover design matters.
High-level architecture
Draw four groups of boxes. First, the app servers. Each runs a cache client library, a small piece of code that knows the list of cache nodes and decides which node owns each key. Some keys that are read extremely often are also kept for a few seconds in the app server's own memory, a local cache. Second, the routing step. It can live inside the client library, or in a proxy: a separate small server that the app talks to, which then forwards each request to the right node. Facebook's mcrouter and Twitter's twemproxy are proxies of this kind. Third, the cache nodes. Each is a machine that keeps keys and values in RAM, evicts old items when full, and expires items when their TTL runs out. Each node may have a replica, a second machine holding a copy. Fourth, the source of truth: the database, plus an invalidator. The invalidator reads the database's commit log (the ordered list of every change the database has saved) and sends a delete to the cache for each key whose data changed. Reads go app, cache, and on a miss the database. Writes go to the database first, and the cached copy is deleted, either by the app or by the invalidator. Cache nodes do not talk to each other in this design, which keeps them simple. Facebook's paper describes its memcached servers this way too.
In a real interview, sketch this on the whiteboard before diving into any single box.
One read and one write, end to end
Follow a GET for a user's profile, then a change to that profile. This is a good order to explain it in an interview.
- 1
GET: local check
The app server needs user:42. If user:42 is a known hot key, it first checks its own local cache, which holds such keys for about 2 seconds.
- 2
GET: pick the node
The client library hashes the key user:42 to a number and finds the first ring point clockwise from it. That point belongs to cache node B.
- 3
GET: hit
Node B finds user:42 in its hash table, moves it to the front of its LRU list, and returns the value. Most reads end here, in well under a millisecond.
- 4
GET: miss
If B does not have it, B says so, and with leases it also hands back a token that allows this one client to refill the key. Other clients that miss now are told to wait a moment.
- 5
GET: refill
The client reads user:42 from the database and calls set on B with the value, a TTL of, say, 300 seconds, and the token. B stores it only if the token is still valid.
- 6
SET: write the database
The user changes their name. The app writes the new name to the database. The database is updated first because it is the source of truth.
- 7
SET: delete the key
The cached user:42 is now out of date, so it is deleted from node B (and from any copies). It is deleted, not overwritten, because a delete can safely be repeated.
- 8
SET: backstop
The invalidator, reading the database's commit log, also sends delete user:42. If the app crashed before its own delete, this one still arrives. The TTL is the last safety net.
- 9
Next GET
The next read misses, refills from the database, and caches the new name.
Core components
Walk through each service. The interviewer wants to hear what each one owns, not just the names.
Cache client library
Code inside each app server. It keeps the list of cache nodes, hashes each key to its node, sends requests, and batches many gets into one round trip. It also marks a node as down after repeated timeouts. Facebook's paper says its client handles serialization (turning objects into bytes), compression, routing, error handling and batching, and that the server list comes from a separate configuration system.
Routing proxy (optional)
A small server between the apps and the cache nodes. Apps talk to it as if it were one cache. It routes each key, keeps a few connections open to each node instead of one per app thread, and can fail over to another node. mcrouter (Facebook) and twemproxy (Twitter) are real examples. The cost is one more network hop.
Cache node
A machine running the cache server. It holds a hash table from key to item, a memory allocator, and an eviction list. It answers get, set, add and delete. It does not know about other nodes and does not replicate on its own in the Memcached model.
Hash ring (consistent hashing)
The map from keys to nodes. Each node is placed at many points on a circle of hash values. A key belongs to the first node point clockwise from the key's own hash. Adding or removing a node moves only the keys next to its points.
Eviction and expiry
Eviction removes items when memory is full; LRU (least recently used) is the common choice, and LFU (least frequently used) is another. Expiry removes items whose TTL has passed. Memcached checks for already expired items at the tail of the LRU list before it evicts a live one.
Replicas and spare pool
Copies that take over when a node dies. A replica is a second machine holding a copy of a node's keys. A spare pool (Facebook calls it Gutter) is a small set of idle machines, about 1% of the cluster, that briefly caches the keys of failed nodes with a short TTL.
Invalidator
A process next to the database that reads its commit log and sends a delete for every cached key whose data changed. Facebook's version is called mcsqueal. It batches deletes into fewer network packets.
Configuration and health service
Holds the current list of cache nodes and their ring points, and tells every client and proxy when it changes. It must change the list carefully: two clients with different lists send one key to two different nodes.
Data model
Pick the right store per table. Justify each choice with the access pattern, not by reflex.
item (inside one cache node)key (string)value (bytes)ttl / expires_atflagslast access time or LRU pointersThe cache does not understand the value; it stores bytes. Keys should name what they hold and carry a version when the format changes, for example user:42:v3, so old and new code never read each other's bytes. In memcached, a TTL above 30 days (2,592,000 seconds) is read as an exact date and time (a Unix timestamp) instead of a number of seconds from now.
hash table (inside one cache node)hash(key) -> itemFinds an item in one step on average. Each item also sits in an eviction list, so the node can find the oldest item without scanning.
ring (shared by every client or proxy)point (32-bit hash)nodeA sorted list of points. With 14 nodes and 160 points each (the count twemproxy's ketama code uses), it holds 2,240 entries, small enough to keep in every client. Lookup is a binary search, which halves the sorted list at each step: about 11 steps for 2,240 points.
lease (inside one cache node, short-lived)keytoken (64-bit)issued_atRecords which client may refill a missing key. Facebook issues at most one token per key every 10 seconds. A delete of the key cancels the token, so a refill that started before the delete is refused.
Deep dives
These are the conversations the interviewer is steering you toward. Practice each one until you can talk through it without notes.
Cache-aside, write-through and write-back
There are three common ways to keep a cache and a database together. They differ in who writes what, and in what can go wrong. In cache-aside, also called lazy loading, the app does everything. On a read, it asks the cache. On a miss, it reads the database and puts the value into the cache. On a write, it updates the database and then deletes the cached key. Only data someone asked for gets cached, and a dead cache node only makes reads slower. The AWS ElastiCache guide lists the cost: a miss takes three trips (cache, database, cache again), and cached data can go stale. Facebook uses this pattern and calls it a demand-filled look-aside cache. Its paper explains why it deletes instead of updating the cached value: deletes are idempotent, which means doing one twice has the effect of doing it once. Two updates that arrive in the wrong order would leave the older value cached. In write-through, every write goes to the cache layer and to the database before the app gets an answer. Cached data is current, but each write takes two steps, and much of what is written is never read. A new, empty node also stays empty until its keys are written again, so in practice write-through is combined with lazy loading. In write-back, also called write-behind, the cache replies at once and writes the database later. Writes are very fast, and many writes to one key can be merged into one. But if the cache node dies before the write reaches the database, data the app was told was saved is gone. A general cache design should start with cache-aside plus a TTL on every key. Use write-back only for data you can afford to lose, such as view counters.


Eviction: LRU, LFU and TTL
Memory is limited, so a full node must throw something away. That is eviction. The goal is to keep the items most likely to be read again. LRU, least recently used, throws away the item that was read longest ago. It works well because popular items tend to be read again soon. The classic way to do it in O(1) time, meaning a fixed amount of work no matter how many items are stored, uses two parts. A hash map finds a key's entry in one step. A doubly linked list keeps entries in order of last use, newest at the front. Each list entry points both forward and backward, so it can be cut out without walking the list. A read moves the entry to the front. A write to a full cache removes the entry at the back. Real systems bend this idea. Redis does not keep an exact list, because the pointers cost memory. It samples a few keys at random, 5 by default, and evicts the oldest of them. Its docs say this comes very close to true LRU. Memcached splits its LRU into HOT, WARM and COLD parts. New items enter HOT, and items that are read again move to WARM. This protects busy items from a one-time scan of many keys that are never read again. LFU, least frequently used, evicts the item read the fewest times. It is better when some keys are steadily popular. It needs a counter that fades over time, or an item that was popular last week never leaves. Redis's LFU uses a small probabilistic counter that is reduced every minute by default. TTL is separate from eviction. A TTL says how long an item may live even if there is room, which bounds how stale it can get. Set a TTL on every key. Add a little randomness to it (for example 300 seconds plus or minus 30), so keys written together do not all expire together.
class Node:
__slots__ = ("key", "val", "prev", "next")
def __init__(self, key=None, val=None):
self.key, self.val, self.prev, self.next = key, val, None, None
class LRUCache:
def __init__(self, capacity):
self.cap, self.map = capacity, {} # map: key -> list node
self.head, self.tail = Node(), Node() # two empty end markers
self.head.next, self.tail.prev = self.tail, self.head
def _unlink(self, n):
n.prev.next, n.next.prev = n.next, n.prev
def _push_front(self, n): # front = most recently used
n.prev, n.next = self.head, self.head.next
self.head.next.prev = n
self.head.next = n
def get(self, key):
n = self.map.get(key)
if n is None:
return None # miss
self._unlink(n)
self._push_front(n)
return n.val
def set(self, key, val):
n = self.map.get(key)
if n is not None:
n.val = val
self._unlink(n)
else:
if len(self.map) >= self.cap: # full: drop the oldest
old = self.tail.prev
self._unlink(old)
del self.map[old.key]
n = Node(key, val)
self.map[key] = n
self._push_front(n)
c = LRUCache(2)
c.set("a", 1); c.set("b", 2)
c.get("a") # a is now the newest
c.set("c", 3) # full: b is the oldest, so b goes
print(c.get("b"), c.get("a"), c.get("c")) # None 1 3
Sharding with consistent hashing and virtual nodes
Sharding means splitting the keys so that each node holds only part of them. The simplest rule is node = hash(key) mod N, where N is the number of nodes and mod is the remainder after division. It spreads keys well, but it breaks badly when N changes. Going from 14 to 15 nodes changes the remainder for about 14 keys in 15. In our test, 93.3% of keys moved. For a cache, a moved key is a miss, so adding one node would empty most of the cache and send that load to the database. Consistent hashing fixes this. It comes from a 1997 paper by Karger and others, written to spread web caching load; the paper describes a hash function that 'changes minimally' when the set of caches changes. Picture a circle of numbers from 0 to 2^32 - 1. Each node is hashed to points on that circle. Each key is hashed to a spot on the circle too, and it belongs to the first node point clockwise. When a node joins, it takes over only the keys just before its points. When a node leaves, only its keys move. With one point per node, though, the gaps between points are uneven, and the circle gives one node a much bigger share. The fix is virtual nodes: each real node gets many points. Amazon's Dynamo paper lists the gains. A failed node's load spreads evenly over the remaining nodes. A new node takes a little load from each existing node. And a bigger machine can be given more points. Twemproxy's ketama code uses 160 points per server. Running the code below on 200,000 keys and 14 nodes: with 1 point per node, the busiest node held 2.7 times its fair share; with 160 points, 1.13 times. Adding a 15th node moved 6.9% of keys, close to the ideal of 1 in 15 (6.7%). Redis Cluster uses a different scheme: a fixed set of 16,384 hash slots, slot = CRC16(key) mod 16384 (CRC16 is a simple checksum function), with each primary owning a range of slots. Moving a slot moves its keys, and clients learn the new owner from a redirect reply.
import bisect, hashlib
def h(s): # a 32-bit spot on the ring
return int(hashlib.md5(s.encode()).hexdigest()[:8], 16)
class Ring:
def __init__(self, nodes, vnodes=160):
self.points = sorted((h(f"{n}#{i}"), n) for n in nodes for i in range(vnodes))
self.spots = [p for p, _ in self.points]
def node_for(self, key):
i = bisect.bisect_right(self.spots, h(key)) # first point clockwise
return self.points[i % len(self.points)][1] # past the top: wrap to 0
keys = [f"user:{i}" for i in range(200_000)]
old = Ring([f"cache-{n}" for n in range(14)])
new = Ring([f"cache-{n}" for n in range(15)])
moved = sum(old.node_for(k) != new.node_for(k) for k in keys) / len(keys)
print(f"{moved:.1%} of keys moved")

Routing: client library, proxy, or redirects
Something has to turn a key into a node address. There are three places to put that logic. The first is a client library inside each app. It is the fastest, because there is no extra hop. Memcached works this way: the servers do not know about each other, and the docs note that finding the server is just a hash lookup on the client. The weak point is that every client must hold an identical node list. Karger's paper calls each client's list its 'view', and its analysis allows views to differ, but two different lists still send one key to two nodes for a while. The second is a proxy. Apps send everything to a local or nearby proxy, which routes the request. Facebook's mcrouter is described in its README as a memcached protocol router, and it handles almost 5 billion requests a second at peak across Facebook and Instagram. It adds connection pooling (keeping a few connections open and reusing them), replicated pools, failover and cold cache warm-up. Twemproxy was built mainly to cut the number of connections to the cache servers. It merges many client connections into a few server connections, and supports ketama consistent hashing. Facebook uses both styles at once: gets go straight from the client to memcached over UDP, a lighter network protocol, while sets and deletes go over TCP, the usual reliable protocol, through an mcrouter running on the web server itself. The third is redirects from the servers themselves. In Redis Cluster, a client may ask any node; a node that does not own the key replies MOVED with the right address, and smart clients cache the slot map so later requests go direct. Pick the client library when you control every app, a proxy when many languages or very many connections are involved, and redirects when you use Redis Cluster. The Design Redis page explains the slot map and resharding in more depth.
Replication and failover
A cache can lose data, because the database still has it. So why replicate at all? Because a dead node does not just lose data; it moves load. In our numbers, losing 1 of 14 nodes with no copy sends database reads from 50,000 to about 118,000 a second. There are three common answers. The first is a replica per node. Each primary (the main copy) sends its writes to a replica, and if the primary dies, the replica is promoted and serves its keys. Replication is usually asynchronous, which means the primary answers before the replica has the write. The Redis Cluster spec says plainly that there are small windows where acknowledged writes can be lost. For a cache this is usually fine, with one trap: a lost delete. If a delete reached the old primary but not the replica, the promoted replica still serves the old value. The TTL is what limits how long that lasts. The second answer is a spare pool. Facebook's Gutter is a small set of machines, about 1% of a cluster. When a client gets no answer from a node, it retries against Gutter, and on a Gutter miss it reads the database and fills Gutter, with a short expiry. Facebook reports this cut client-visible failures by 99% and turned 10% to 25% of failures into hits each day. The third answer is the tempting wrong one: rehash the dead node's keys onto the remaining nodes. Facebook's paper warns against it, because one key can account for 20% of a server's requests, and the node that inherits a hot key can be overloaded in turn. Failure detection matters too. Clients mark a node down after a few timeouts in a row; twemproxy has a setting for this, auto_eject_hosts. Mark it down too quickly and a short network blip moves keys for nothing. Too slowly, and requests wait on a dead machine. Finally, an empty cluster is dangerous. Facebook's cold cluster warmup lets a new cluster read from a warm one instead of the database, and says this brings a cold cluster to full capacity in a few hours instead of a few days.

Hot keys
Consistent hashing spreads keys evenly, but it cannot spread a single key. If one celebrity's profile is read 200,000 times a second, the one node that owns it gets all 200,000 on top of its normal load. In our numbers that node goes from 71,400 to 271,400 reads a second, far over a 100,000 budget. You first have to find hot keys: count reads per key on each node or in the client, and alert on keys above a threshold. Then there are three fixes. The first is a local cache in each app server, holding hot keys for a second or two. With 1,000 app servers and a 2-second local TTL, at most 1,000 / 2 = 500 reads a second for that key reach the cache cluster. The cost is that each app server may be up to 2 seconds behind. The second is to store several copies of the key under different names, such as key#0 to key#7, which land on different nodes. Readers pick one at random, so each copy gets 200,000 / 8 = 25,000 reads a second. Every write must delete all eight copies. The third is replication of a whole group of keys inside a pool. Facebook does this when an app reads many keys at once, the whole set fits on one or two servers, and the request rate is more than one server can manage. Its paper shows why splitting the keys would not help in that case: a request for 100 keys split over two servers still costs each server one request, so each server's request rate does not drop. With full copies, each request goes to one replica, and the load per server halves. One more source of load is reads for rows that do not exist. If a missing user id is never cached, every request for it reaches the database, so cache the 'not found' answer too, with a short TTL. Hot writes are a different problem: if one key is written thousands of times a second, its cached copy is deleted constantly and every reader misses. That case needs the stampede protection in the next section, or a cache with a short TTL that is not deleted on every write.

Cache stampedes: leases, request coalescing and early refresh
A cache stampede, also called a thundering herd, happens when a popular key disappears, by expiry, eviction or a delete, and many requests miss at once. Every one of them goes to the database to rebuild one value. The database slows down, the rebuild takes longer, and more requests pile up behind it. Facebook's paper describes the cause plainly: a key with heavy reads and writes, where each write invalidates the value and many reads fall through to the costly path. There are three standard fixes. The first is leases, from that paper. On a miss, the cache gives the client a lease, a 64-bit token tied to that key. Only a set that carries a valid token is stored. The server hands out at most one token per key every 10 seconds; other clients that miss in that time are told to wait a short time and retry. Usually the lease holder has set the value within a few milliseconds, so the retries hit. Over one week, for keys prone to stampedes, peak database reads fell from 17,000 to 1,300 a second. Leases also fix stale sets: if a delete for the key arrives while the lease holder is still reading the database, the token is cancelled, so an old value cannot be written over the newer state. The second fix is request coalescing inside each app server: when many threads miss on one key, one of them reads the database and the rest wait for its result. Go's singleflight package does this, making sure that only one call is in flight for a given key at a time. It limits the herd to one request per app server, not one in total. Outside Facebook's memcache, a short lock key does the job of a lease, as in the code below. The third fix refreshes a key before it expires. The XFetch paper (VLDB 2015) has each reader, near the expiry time, decide at random whether to recompute early, with the chance rising as expiry gets closer and as the recompute gets slower. One reader usually refreshes the value while the others keep getting hits. Random jitter on TTLs helps too, by spreading expiries out.
import time
def get_or_load(cache, key, load, ttl=300, lock_ttl=10):
val = cache.get(key)
if val is not None:
return val # hit
if cache.add(f"lock:{key}", 1, lock_ttl): # only one caller wins this
try:
val = load(key) # the one database read
cache.set(key, val, ttl)
finally:
cache.delete(f"lock:{key}")
return val
for _ in range(50): # everyone else: wait up to ~1 s
time.sleep(0.02)
val = cache.get(key)
if val is not None:
return val
return load(key) # still nothing: read the source
Consistency and invalidation
A cache holds copies, and copies can be out of date. The design question is how out of date, and for how long. The most common bug is a race in cache-aside. Reader A misses and reads the old value from the database. Before A writes it to the cache, writer B updates the database and deletes the key. Then A sets the old value. Now the cache holds stale data until its TTL ends. Leases stop this, as described above, because B's delete cancels A's token. A versioned set does too: memcached's cas command stores a value only if no one else has updated the key since the caller last read it. Second, deletes can get lost: the app can crash between writing the database and deleting the key. The fix is to drive deletes from the database itself. Facebook runs a background program called mcsqueal on every database. It reads the statements the database committed, pulls out the cache keys to delete, and sends them to the cache in every front-end cluster (a group of web servers and their cache servers) in that region. It batches them into fewer packets through mcrouter; Facebook measured 18 times more deletes per packet (the median) this way. Most of these deletes find nothing to delete: only 4% of them removed cached data. Third, across regions, a delete can arrive before the database change itself reaches that region's replica database. A read in between would put old data back into the cache. Facebook avoids part of this by sending invalidations from the database side, after the change is committed, and uses markers to send some reads to the main region while a change is still on its way. Last, the TTL is the safety net for every case above. Pick it from how stale the data may be. A product price might allow 60 seconds; a user's own just-saved profile should be deleted on write, never left to the TTL. Facebook's paper sums up its own choice as best-effort eventual consistency, with the emphasis on performance and availability. In an interview, say your own version out loud: which reads may be stale, for how long, and what bounds it.
Memory, multi-gets and the network
Three practical details separate a working cache from a slow one. First, memory is handed out in chunks. Memcached splits its memory into 1 MB pages, gives each page to a slab class (a group of chunks of one size), and cuts it into equal chunks. An item goes into the smallest chunk that fits, so the rest of the chunk is wasted. The docs show a 50-byte item losing 30 bytes in a class of small chunks. Facebook's version uses slab classes from 64 bytes up to 1 MB, each about 1.07 times the one before. That is why the estimate above adds 25% headroom instead of using the raw 550 GB. Second, pages need many keys. On Facebook, loading a popular page fetched 521 distinct items from memcache on average. So the client sends gets in batches, a multi-get, averaging 24 keys per request in Facebook's paper. Third, the network can choke. When one web server asks many cache servers at once, all the replies arrive together and overflow the network switch buffers. This is called incast congestion. Facebook's clients limit how many requests are outstanding with a sliding window, which grows slowly while requests succeed and shrinks when a request goes unanswered. Facebook also sends gets over UDP, which skips the setup cost of a TCP connection; under peak load, 0.25% of those gets were dropped, and the client treats a dropped get as a miss.
Trade-offs to discuss
Every senior interviewer expects you to surface at least 3 of these. Pick the decisions, state the alternatives, and justify your choice.
Cache-aside vs write-through vs write-back
Cache-aside is simple and survives node loss, but the first read after a change is a miss and stale data is possible until a delete or the TTL. Write-through keeps cached data current, at the cost of slower writes and caching data nobody reads. Write-back gives the fastest writes and can lose acknowledged data on a crash. Default to cache-aside with a TTL on every key.
Delete on write vs update on write
Updating the cached value saves the next miss, but two updates that arrive out of order leave the older one cached. Deleting is idempotent and safe to repeat, so Facebook deletes. Pair delete with leases to stop a slow reader from writing an old value back.
Client-side routing vs a proxy
A client library is one network hop shorter and has no extra server to run, but every app in every language needs it, and all of them must agree on the node list. A proxy such as mcrouter or twemproxy centralises routing, pools connections and handles failover, at the cost of one more hop and one more service to keep running.
Replicas vs a spare pool vs no protection
Replicas double the memory bill but keep the hit rate when a node dies. A spare pool like Facebook's Gutter costs about 1% extra machines and fills up within minutes. No protection is cheapest and is fine only if the database can take the extra load of a dead node, about 2.4 times normal in our numbers.
LRU vs LFU
LRU is simple and adapts at once when interest moves to new keys, but one scan over many cold keys can push out the hot ones. LFU keeps steadily popular keys, but needs counters that fade over time, or yesterday's favourites never leave. Segmented LRU, as in memcached's HOT, WARM and COLD lists, sits in between.
Long TTL vs short TTL
A long TTL means a higher hit rate and less database load, and stale data lasts longer if a delete is lost. A short TTL bounds staleness but adds misses and makes stampedes more likely. Choose per type of data, and add random jitter so keys do not expire together.
Local in-process cache vs shared cache only
A small local cache absorbs hot keys and saves a network trip, but each app server can be out of date for its local TTL, and deletes do not reach it. Keep local TTLs to a second or two, and only for keys that tolerate it.
Free PDF · 18 pages
Get the free System Design Interview Cheat Sheet
The interview in seven stages, the numbers worth knowing by heart, and twelve classic systems on one page each, every line linked to the lesson it comes from.
Follow-up questions to expect
After the main design, the interviewer usually picks one part and asks more. Here is what to say.
You add a 15th node. What happens to the hit rate?
With consistent hashing and virtual nodes, about 1 key in 15 moves to the new node, so the hit rate dips by a few percent and recovers as those keys are read again. With hash mod N, about 93% of keys would move, and the database would take most of the read load at once.
A cache node dies at peak. Walk me through the next minute.
Clients see timeouts and mark the node down. Its keys go to a promoted replica or to a spare pool, depending on the design. Without either, about 71,400 reads a second become database reads. Watch for the hot keys that lived on that node, and do not rehash them onto one neighbour.
How would you find hot keys in production?
Count requests per key, either in the client or proxy or by sampling on the cache nodes, and keep the top keys per minute. Alert when one key passes a share of a node's budget. Then serve those keys from a short local cache or from several copies.
Why delete the cached key on a write instead of setting the new value?
A delete can be repeated safely. Two sets that arrive in the wrong order leave the older value in the cache. Facebook's paper gives this reason for deleting.
The cache shows a user their old name right after they changed it. Why?
Most likely the cache-aside race: a reader fetched the old row before the update, then cached it after the delete. Fix it with leases or a versioned set (cas), and drive deletes from the database's commit log so a crashed app cannot skip one.
How is this different from designing Redis?
Designing Redis is about one server: its data types, its single thread and its persistence to disk. This design is about the cluster around many cache servers: how keys are spread, how clients route, what happens on failure, and how the cache stays close to the database.
What would you monitor?
Hit rate per node and overall, evictions, expirations, latency at the 99th percentile (the time 99% of requests beat), requests per node to spot hot keys and uneven hashing, and database reads a second. Redis's INFO command, for example, reports keyspace_hits, keyspace_misses and evicted_keys. A falling hit rate with many evictions means the cache is too small or is evicting the wrong keys.
How do you warm an empty cache?
Do not point full traffic at it. Let it read misses from a warm cache instead of the database, as Facebook's cold cluster warmup does, or move traffic over gradually. Facebook reports a few hours to full capacity this way instead of a few days.
How a Distributed Cache actually does it
The best documented distributed cache is Facebook's memcache, described in 'Scaling Memcache at Facebook' (NSDI 2013). It handled billions of requests a second and held trillions of items. It is a demand-filled look-aside cache: web servers read memcache first, fill it on a miss, and delete keys after writing to MySQL. Keys are spread over hundreds of memcached servers per cluster by consistent hashing, and the servers never talk to each other. Leases solve two problems: stale sets and thundering herds. One token per key every 10 seconds cut peak database reads for stampede-prone keys from 17,000 to 1,300 a second. A Gutter pool of about 1% of the servers absorbs the load of failed ones, cutting client-visible failures by 99%. An invalidation daemon, mcsqueal, reads the database's commits and sends deletes in batches. Facebook's mcrouter, now open source, routes memcached requests and handles almost 5 billion requests a second at peak across Facebook and Instagram. Twitter's twemproxy is a lighter proxy for memcached and Redis, built to cut connections to the cache servers, with ketama consistent hashing at 160 points per server. The ring itself comes from Karger and others (STOC 1997), and Amazon's Dynamo paper (SOSP 2007) made virtual nodes standard. Redis Cluster splits keys into 16,384 slots, uses asynchronous replication, and says openly that acknowledged writes can be lost in some failovers. Memcached's own docs describe client-side hashing, the slab allocator, and a segmented HOT, WARM and COLD LRU. Two research results finish the picture: the XFetch paper on refreshing keys early at random to prevent stampedes, and Go's singleflight package for collapsing duplicate misses inside one process.
Sources
- Scaling Memcache at Facebook (NSDI 2013)
- mcrouter, Facebook's memcached protocol router (README)
- twemproxy (nutcracker), Twitter's memcached and Redis proxy (README)
- twemproxy ketama code: 160 points per server
- Memcached docs: performance, slabs and client-side hashing
- Memcached: the segmented HOT, WARM and COLD LRU
- Memcached protocol: set, add, cas and expiry times
- Redis Cluster specification (hash slots, redirects, write safety)
- Redis docs: key eviction (approximated LRU, LFU)
- Karger et al., Consistent Hashing and Random Trees (STOC 1997)
- Dynamo: Amazon's Highly Available Key-value Store (SOSP 2007), virtual nodes
- Optimal Probabilistic Cache Stampede Prevention, XFetch (VLDB 2015)
- Go singleflight: duplicate call suppression
- AWS ElastiCache: caching strategies (lazy loading, write-through, TTL)
- Wikipedia: Cache (computing), writing policies
Lessons to study before this interview
If any of these topics are fuzzy, the interviewer will catch it. Each lesson is 15 to 60 minutes with diagrams, code, and a quiz.
Distributed Cache
foundation / caching strategies
Consistent Hashing
intermediate / data replication distribution
Cache Eviction Policies
foundation / caching strategies
Cache-Aside Pattern
foundation / caching strategies
Write-Through Cache
foundation / caching strategies
Write-Back Cache
foundation / caching strategies
Cache Stampede Prevention
foundation / caching strategies
Cache Invalidation
foundation / caching strategies
Time To Live (TTL)
foundation / caching strategies
Cache Warming
foundation / caching strategies
Memcached
foundation / caching strategies
Asynchronous Replication
intermediate / data replication distribution
Failover
advanced / reliability resilience
Frequently asked questions
Practice with 770 system design lessons
Lifetime access for ₹499 in India or $49 elsewhere. Interactive diagrams, quizzes, and 20 capstone projects to practise on.
