Design a Distributed Cache
A distributed cache is a fast, in-memory key-value store spread across many machines, used to keep hot data close to the application and shield slower stores. Almost every other design in this track uses one, so interviewers like to ask you to build one. The question tests whether you understand partitioning, replication, eviction and failure from first principles, and it rewards a clear explanation of consistent hashing.
The chapter follows the usual shape: understand the problem, set up the interface, build the high-level design, then go deep on the questions interviewers use to separate levels.
1. Understanding the problem
We are building a service that stores key-value pairs in memory across a cluster, with very fast reads and writes, and that keeps working when machines fail or when capacity is added.
Functional requirements
Core:
- Get, set and delete a value by key, with an optional time-to-live.
- Data is spread across many nodes, so the cache is larger than any one machine.
- When memory is full, the cache evicts entries according to a policy.
Confirm in or out: richer data types such as lists and counters, atomic operations, persistence to disk, transactions, and pub-sub. A sensible opening: "I will design a partitioned, replicated key-value cache with TTLs and eviction, and treat rich data types and persistence as extensions."
Non-functional requirements
- Very low latency. Reads in well under a millisecond on the server, a millisecond or two with the network.
- High throughput. Hundreds of thousands of operations per second per node.
- High availability. A cache outage pushes full load onto the database, which can then fail, so the cache must tolerate node failures.
- Scalability. Add and remove nodes with minimal disruption.
- Acceptable consistency. A cache is allowed to be a little stale or to lose an entry, because the source of truth lives elsewhere.
Estimation
Assume the application needs to cache 2 TB of hot data and 1 million operations per second.
| Quantity | Calculation | Result |
|---|---|---|
| Memory per node | say 64 GB usable | |
| Nodes for capacity | 2 TB / 64 GB | about 32 primaries |
| With one replica each | about 64 nodes | |
| Operations per node | about 31,000 per second, comfortable |
What the numbers say. Capacity, not throughput, drives the node count. Each node handles an easy load, which is why the design effort goes into placement, failure and eviction rather than raw speed.
2. The set up
Core entities
- Entry: key, value, expiry time, and bookkeeping for eviction.
- Node: a server holding a subset of the keys.
- Shard (partition): the set of keys one node, with its replicas, owns.
API
GET /cache/{key} -> value or miss
PUT /cache/{key} body, ttl=seconds -> ok
DELETE /cache/{key} -> ok
In practice the protocol is a compact binary or text protocol over persistent connections, not HTTP, but the operations are the same.
3. High-level design
Three questions define the design: where does each key live, how do clients find it, and what happens when a node dies?
- Partitioning. Keys are divided among nodes by hashing the key, so each node owns a share of the key space.
- Routing. Either a client library knows the cluster map and talks directly to the right node, or a proxy layer routes requests. Client-side routing saves a network hop, while a proxy keeps clients simple and lets you change topology without redeploying applications.
- Replication. Each shard has a primary and one or more replicas on other machines, so a node failure does not lose the shard.
A common way to define partitions is a fixed number of hash slots, such as 16,384. A key's slot is a hash of the key modulo the slot count, and each node owns a range of slots. The mapping from slot to node is the cluster map. Moving a slot between nodes is the unit of rebalancing.
4. Potential deep dives
Deep dive 1: How do you assign keys to nodes?
The challenge. Many nodes, billions of keys, and nodes that come and go. A key must be found on the same node every time, and adding capacity should not wreck the cache.
Weak: hash the key and take it modulo the number of nodes. It is simple, and it breaks when changes. Going from 10 to 11 nodes remaps almost every key, since the hash modulo changes for nearly all of them. The cache suddenly misses on nearly everything, and the database takes the entire load at once. This is the classic cache avalanche on a resize.
Solid: consistent hashing. Place the nodes on a hash ring. A key belongs to the first node clockwise from the key's hash position. When a node is added or removed, only the keys between it and its neighbour move, about of the total, so the cache stays warm.
Excellent: consistent hashing with virtual nodes, or fixed hash slots. With few nodes on the ring, the arcs between them are uneven, so some nodes own far more keys than others. Give each physical node many points on the ring, called virtual nodes. This evens out the distribution, lets a larger machine take more points, and spreads the load of a failed node across many survivors instead of dumping it on one neighbour.
<!--fig:ring-->An alternative with the same goals is a fixed set of hash slots assigned to nodes, as described above. A fixed slot table is simple to reason about and easy to rebalance by moving whole slots, while a ring needs no central table but needs care in how it is distributed. Either is a good answer if you explain the properties: balance, minimal movement on change, and a lookup that every client computes the same way.
Deep dive 2: What happens when a node fails?
The challenge. A node holds a share of the data. When it dies, requests for its keys miss, and the database sees a spike.
Weak: do nothing, and let the keys be reloaded on demand. The share of the cache on the failed node starts cold. Misses concentrate on the database, and a popular key can trigger a stampede.
Solid: replicate each shard. Every primary streams its writes to one or more replicas. When the primary dies, a replica is promoted. Replication is usually asynchronous, which is fast but can lose the last few writes at failover. For a cache that is acceptable, because the data can be reloaded from the source of truth.
Excellent: automatic failover with failure detection and a safe rebuild. Nodes exchange heartbeats. When a majority of nodes agree a primary is unreachable, a replica is promoted and the cluster map is updated and propagated to clients. Guard against split brain, where a primary that is merely partitioned continues accepting writes, by requiring a majority to agree and by making a demoted primary stop serving writes. After a failure, bring a new replica up and fill it by copying from the primary, throttled so it does not disturb live traffic. Tell clients about map changes through the cluster itself, with redirects when a client's map is stale, so that routing recovers without manual action.
Mention the trade-off between an unavailable key and a stale one during failover, and say you would accept brief unavailability of a shard rather than serving inconsistent data for a cache that can simply be reloaded.
Deep dive 3: How do you decide what to evict?
The challenge. Memory is finite. When it fills, something must go, and the choice affects the hit rate directly.
Weak: random eviction, or first in first out. Random eviction ignores usefulness. First in, first out removes old entries regardless of how popular they are.
Solid: least recently used (LRU). Remove the entry that has gone longest without being accessed, on the assumption that recent use predicts future use. Implement it with a hash map for lookup and a doubly linked list ordering entries by recency, so get, set and eviction are all constant time. Every access moves the entry to the front of the list, and eviction removes the tail.
Excellent: approximate policies, frequency awareness and expiry. An exact LRU list needs a lock and extra pointers per entry, which costs memory and hurts concurrency. Production caches often use an approximate LRU: sample a few random entries and evict the least recently used among them, which is nearly as effective and far cheaper. Least frequently used (LFU) keeps entries that are popular over time, protecting them from being flushed by a one-off scan, a problem that plain LRU suffers from. Combine eviction with time-to-live: expired entries are removed lazily on access and by a background sweep that samples keys, so memory is reclaimed without scanning everything. Let the application choose the policy per cache, because the right one depends on the access pattern. Track the hit rate and evictions per second, since they tell you whether the cache is too small.
Deep dive 4: Hot keys and stampedes
The challenge. Access is skewed. One viral item can receive more traffic than a single node can serve, and a popular entry expiring can send thousands of requests to the database at once.
Hot key. Weak: ignore it, and the owning node saturates. Solid: detect hot keys and replicate them to several nodes, spreading reads across the copies. Excellent: add a small local cache inside each application server for the hottest keys, with a very short TTL, so most reads never leave the process. Detect hot keys by sampling access counts, and for known hot keys use key suffixes to spread them across shards, for example reading from one of several copies chosen at random.
Stampede (thundering herd). Weak: each miss goes to the database independently. Solid: request coalescing: the first request that misses fetches the value while the others wait for its result, so the database sees one query. Excellent: the same, plus refreshing before expiry with jittered TTLs so popular entries rarely expire under load, and serving a stale value while a background refresh happens. Also cache negative results (a short-lived "not found") so requests for missing keys do not all hit the database, and consider a Bloom filter for keys that definitely do not exist.
Deep dive 5: Consistency with the source of truth
The challenge. The database changes, and the cache must not serve wrong data for long.
Weak: write the database and hope the cache catches up. Stale entries live until their TTL expires.
Solid: cache-aside with invalidation. Reads go to the cache first and fill it on a miss. On a write, update the database and then delete the cache entry, so the next read reloads it. Deleting is safer than updating the cached value, because it avoids races where two writers update the cache out of order.
Excellent: understand the races, and bound the staleness. There is a well-known race: a reader misses, reads the old value from the database, a writer updates the database and deletes the cache entry, and then the reader stores its stale value in the cache. The cache now holds old data until the TTL expires. Mitigate it with a short TTL as a safety net, with versioned values, and with leases that let only the latest reader fill the cache. For data that cannot tolerate staleness, do not cache it, or read from the primary. State that a cache gives you eventual consistency bounded by the TTL, and that choosing the TTL is choosing how wrong you can afford to be.
Write strategies. Write-through (write to cache and database together), write-around (write to the database only) and write-back (write to the cache and flush later, which risks loss on a crash) each fit different needs. Explain the trade-offs, and note that write-back suits data where losing the latest writes is acceptable.
Deep dive 6: Memory efficiency, persistence and operations
Memory. Overhead per entry matters at billions of keys. Use compact encodings for small values, avoid per-entry allocation overhead with slab allocators, and compress large values if CPU allows.
Persistence. A cache is usually allowed to lose data, but a cold restart of a large cache floods the database. Some systems periodically snapshot memory to disk or log writes so that a restarted node reloads quickly. Say you would use it for the large caches where warm-up cost is high.
Operations. Monitor hit rate, latency percentiles, memory use, evictions, connection counts and replication lag. Plan capacity from the hit-rate curve, resize with minimal key movement, and test failover regularly. Protect against cache penetration and abusive access patterns with limits.
5. What is expected at each level
Mid-level. You describe a hash map across nodes, mention hashing to choose the node, and propose LRU eviction with its data structures. You mention that caches are used with a database.
Senior. You explain why modulo hashing fails and how consistent hashing, with virtual nodes, fixes it. You replicate shards and describe failover, use cache-aside with invalidation, and handle stampedes and hot keys.
Staff. You discuss the consistency races and how you bound staleness, the cost of cold starts and the strategy for warm-up, client versus proxy routing, multi-region caching and its invalidation challenge, and how you would operate and observe the cache, including how to size it from hit-rate data.
6. Interview questions and model answers
Q: Why not use hash modulo the number of nodes? Changing the number of nodes remaps almost every key, so the cache goes cold all at once and the database is flooded. Consistent hashing remaps only about one over N of the keys.
Q: Why virtual nodes? They even out the key distribution across nodes, let stronger machines own more points, and spread a failed node's load over many survivors.
Q: How do you implement LRU in constant time? A hash map from key to a node in a doubly linked list that orders entries by recency. Access moves the node to the front, and eviction removes the tail.
Q: What happens on a cache node failure? Its replica is promoted after a majority agrees the primary is down, the cluster map is updated, and clients are redirected. Asynchronous replication may lose the last few writes, which is acceptable for a cache.
Q: How do you prevent a stampede on a popular key? Coalesce misses so one request refills the entry, jitter the TTLs, refresh before expiry and serve a slightly stale value during refresh.
Q: Should a write update or delete the cached value? Delete. Updating the cache from two concurrent writers can leave an older value in place. Deleting forces the next read to load the current value, and a short TTL bounds any remaining race.
7. Common mistakes
- Hash modulo N, which causes an avalanche of misses on any resize.
- A cache with no replication that turns each node failure into a database spike.
- Updating cache entries on write instead of deleting them.
- No protection against stampedes or hot keys.
- Unlimited TTLs, so stale data lives forever.
- Treating the cache as the only copy of important data.
- Ignoring the hit rate, the number that tells you whether the cache is working.