A distributed cache system design splits an in-memory key-value store across many nodes using consistent hashing, replicates each partition for failover, and evicts data with a policy such as LRU or LFU when memory fills. In an interview, a strong answer also explains how the application reads and writes through the cache (cache-aside, write-through, or write-back) and how it survives hot keys and thundering herds. This guide walks through that answer in the order an interviewer expects it.
The question is asked directly ("design Redis") or hidden inside a bigger prompt, such as the cache layer of a URL shortener or a news feed.
Key Takeaways
- Start with requirements: data size, QPS, read/write ratio, latency target, and whether stale reads are acceptable. These numbers decide everything else.
- Use consistent hashing with virtual nodes (or a fixed slot map, like Redis Cluster's 16,384 slots) so adding or losing a node moves only a fraction of keys.
- Replicate each shard asynchronously with automatic failover, and say out loud that async replication can lose recent writes.
- Pick an eviction policy from the workload, not from habit: LRU for recency, LFU for stable popularity, TTL for freshness.
- Choose a caching pattern by its failure mode. Cache-aside fails toward stale or missing data, write-back fails toward lost writes.
- Name hot keys and thundering herds before the interviewer does, and offer request coalescing, local caching, and TTL jitter as fixes.
What Are the Requirements and API for a Distributed Cache?
The requirements step sets the scope. Ask questions, then write down what you agreed on.
Functional requirements
get(key)returns the value or a miss.set(key, value, ttl)stores a value with an optional expiry.delete(key)removes a value, mostly for invalidation.- Values are opaque bytes, typically under 1 MB.
Non-functional requirements
- Low latency: single-digit millisecond p99 inside a data center.
- High availability: losing one node should not take the cache down.
- Horizontal scale: add nodes to grow capacity and throughput.
- Eventual consistency with the database is acceptable for most keys.
A minimal API looks like this:
interface CacheClient {
get(key: string): Promise<Uint8Array | null>;
set(key: string, value: Uint8Array, ttlSeconds?: number): Promise<void>;
delete(key: string): Promise<void>;
mget(keys: string[]): Promise<Map<string, Uint8Array>>;
}
Then do quick capacity math out loud. Suppose 2 TB of hot data, 1 million reads per second, and a 10:1 read/write ratio. If each node holds about 64 GB of usable memory after overhead, you need around 32 primary shards for data alone, plus replicas. If one node can serve a few hundred thousand simple operations per second, throughput is not the bottleneck at that size; memory is. Treat per-node throughput as a stated assumption, since it depends on value size, pipelining, and hardware.
How Do You Partition Data with Consistent Hashing?
Partitioning decides which node owns each key, and consistent hashing is the standard answer because it keeps data movement small when the cluster changes.
Start with the naive approach so you can reject it. With node = hash(key) % N, going from 10 to 11 nodes remaps roughly 10/11 of all keys. Every one of those becomes a cache miss at once, and the database takes the full load.
Consistent hashing fixes this. Hash both nodes and keys onto a ring of values from 0 to 2^32 - 1. Each key belongs to the first node found walking clockwise from the key's position.
0 / 2^32
|
C#2 .-----+-----. A#1
/ \
k3 → | | ← k1 (k1 → A#1)
| hash ring |
B#2 →| |← B#1
\ /
A#2 '-----+-----' C#1
|
k2 (→ C#1)
A#1, A#2 = virtual nodes for server A
When server D joins, it takes over only the arcs just counter-clockwise of its positions. About 1/N of the keys move, and the rest stay put.
Virtual nodes make this work in practice. With one point per server, arcs are uneven and one server might own 40% of the ring. Giving each server 100 to 200 virtual points evens out load. It also means a failed server's keys scatter across all survivors instead of landing on one neighbor, which would otherwise double that neighbor's load.
import bisect
import hashlib
class HashRing:
def __init__(self, nodes, vnodes=150):
self.vnodes = vnodes
self.ring = []
self.owners = {}
for node in nodes:
self.add(node)
def _hash(self, value):
return int(hashlib.md5(value.encode()).hexdigest()[:8], 16)
def add(self, node):
for i in range(self.vnodes):
h = self._hash(f"{node}#{i}")
bisect.insort(self.ring, h)
self.owners[h] = node
def get(self, key):
h = self._hash(key)
idx = bisect.bisect(self.ring, h) % len(self.ring)
return self.owners[self.ring[idx]]
The slot alternative. Redis Cluster does not use a ring. According to the Redis cluster specification, the key space is split into 16,384 hash slots, and each key maps to CRC16(key) mod 16384. Each primary owns a set of slots, and rebalancing means migrating whole slots between nodes. Hash tags, the part of a key inside {}, force related keys into the same slot so multi-key operations still work. Mentioning that both designs exist, and that slots make ownership explicit and easy to gossip, is a good senior-level signal.
Where routing lives is the last decision. You have three options:
| Routing option | How it works | Pros | Cons |
|---|---|---|---|
| Smart client | Client library holds the ring or slot map | No extra hop, lowest latency | Every client language needs the logic; map updates must reach all clients |
| Proxy tier | Clients talk to a proxy (Twemproxy or Envoy style) that routes | Thin clients, central control | Extra network hop, proxy fleet to run |
| Server redirect | Any node accepts the request and redirects if it is not the owner | Simple clients can bootstrap | Extra round trip on stale maps |
Redis Cluster combines smart clients with redirects: a client that hits the wrong node gets a MOVED reply and updates its map.
How Do Replication and Failover Work?
Replication keeps a copy of each shard on another machine so a node failure does not erase a slice of the cache. A typical setup gives each primary one or two replicas in different racks or availability zones.
Writes go to the primary, which streams them to replicas asynchronously. That keeps write latency low, but it has a cost you should state plainly: if a primary acknowledges a write and dies before the replica receives it, that write is gone. For a cache this is usually fine because the database is the source of truth. If the interviewer pushes on it, you can mention synchronous acknowledgment from one replica (Redis offers a WAIT command for this), while noting it narrows the loss window rather than giving strong consistency.
Failover needs failure detection and a decision process:
- Nodes exchange heartbeats (gossip in Redis Cluster, or a coordinator such as ZooKeeper or etcd in custom designs).
- If a primary misses heartbeats past a timeout, peers mark it as suspected, then failed once a majority agrees.
- A replica is promoted, and the slot or ring map is updated with a new version number (an epoch).
- Clients learn the new map through redirects or a config push.
Two failure scenarios are worth raising before you are asked. A network partition can leave an old primary still accepting writes on the minority side; requiring a majority to promote, and having the isolated primary stop accepting writes after it loses contact, limits this split-brain window. Facebook's Scaling Memcache at Facebook paper (NSDI 2013) describes another option for short failures: a small "gutter" pool of spare servers that temporarily absorbs a failed server's traffic with short TTLs, instead of rehashing that load onto other primaries.
Practicing this design is one thing; explaining quorum failover clearly while an interviewer watches is another. TechScreen is an invisible AI interview assistant that gives you real-time prompts during system design rounds on Zoom, Meet, or Teams. Start with 3 free tokens, no credit card.
Which Cache Eviction Policies Should You Compare?
An eviction policy decides which key to drop when memory is full, and the right one depends on your access pattern. Interviewers want to hear the tradeoff, not a list of acronyms.
The Redis eviction documentation lists the options a real system exposes: noeviction, allkeys-lru, allkeys-lfu, allkeys-random, volatile-lru, volatile-lfu, volatile-random, and volatile-ttl. The volatile-* policies only consider keys that have a TTL. Redis approximates LRU and LFU by sampling a few keys instead of keeping a perfect ordering, which saves memory.
Use this decision table in the interview:
| Workload signal | Policy | Why it fits | Failure it prevents | Watch out for |
|---|---|---|---|---|
| General web reads, recency matters | LRU (allkeys) | Recently used keys tend to be used again | Stale long-tail keys hogging memory | A single large scan can flush the hot set |
| Stable popular set plus one-off scans | LFU | Frequency protects popular keys from scan bursts | Cache pollution from batch jobs | New hot keys need time to build counts; counters must decay |
| Data with natural freshness windows | TTL / volatile-ttl | Evicts what would expire soonest anyway | Serving data past its useful life | Keys with no TTL are never candidates |
| Mixed: some keys must persist, others are disposable | volatile-lru or volatile-lfu | Only keys with a TTL are eligible | Losing configuration or session keys | Behaves like noeviction if nothing has a TTL |
| Uniform access, no clear pattern | Random | Nearly free to compute | Overhead of tracking recency | Can drop hot keys by chance |
| Cache is really a primary store | noeviction | Writes fail instead of silently dropping data | Silent data loss | Writes error out at full memory |
A good follow-up is how LRU works at the data structure level: a hash map pointing into a doubly linked list gives O(1) get, insert, and move-to-front. Memcached uses a segmented LRU (hot, warm, and cold queues) per slab class, which is worth one sentence if you pick Memcached as your reference.
Which Caching Pattern Should You Use, and How Do You Handle Invalidation?
A caching pattern defines how reads and writes flow between the application, the cache, and the database. Each pattern has a distinct failure mode, and choosing one is really choosing which failure you can live with.
| Pattern | Read path | Write path | Consistency | Failure scenario | Best for |
|---|---|---|---|---|---|
| Cache-aside (lazy loading) | App reads cache; on miss reads DB and fills cache | App writes DB, then deletes the cache key | Eventual | Race: a slow reader fills an old value after a delete, leaving stale data until TTL | Read-heavy general workloads |
| Read-through | Cache library loads from DB on miss | Usually paired with write-through | Eventual | Cache tier outage blocks reads unless the app can fall back | Teams that want loading logic in one place |
| Write-through | Reads hit a warm cache | Write to cache and DB synchronously | Stronger | Write succeeds in DB but fails in cache, or the reverse, without retries | Data read soon after it is written |
| Write-back (write-behind) | Reads hit the cache | Write to cache; flush to DB asynchronously | Weak until flush | Node dies before flush and acknowledged writes are lost | Write-heavy counters, metrics, tolerant data |
| Write-around | Reads fill on miss | Write to DB only | Eventual | First read after a write is always a miss | Data written often but read rarely |
For most interview prompts, say cache-aside with delete-on-write and a TTL as a safety net. Then cover invalidation explicitly.
Cache invalidation strategies, ranked by effort:
- TTL only. Every key expires. Simple, but staleness is bounded only by the TTL.
- Delete on write. The writer deletes the key after committing to the database. Delete rather than update, so two concurrent writers cannot leave the older value behind.
- Change data capture. A process tails the database's change log and issues deletes. This catches writes from every service, including ones that forget to invalidate.
- Versioned keys. Include a version in the key (
user:42:v7), so a write bumps the version and old entries simply age out.
The classic cache-aside race goes like this: reader A misses and reads the old row, writer B updates the row and deletes the key, then reader A writes the old value into the cache. Fixes include short TTLs, a delayed second delete, or leases. Facebook's memcache paper describes leases that reject a fill if the key was deleted after the lease was issued, which closes this exact gap.
How Do You Handle Hot Keys and the Thundering Herd?
Hot keys and thundering herds are the two problems that break a cache which looks fine on the whiteboard. Bring them up yourself.
A hot key is a single key that receives a large share of traffic, such as a celebrity profile or a flash-sale item. Consistent hashing sends all of that traffic to one shard, and adding nodes does not help. Fixes:
- Local in-process cache on each app server with a short TTL (one to a few seconds). This absorbs most of the load before it reaches the network.
- Key replication with suffixes. Store copies as
item:9:0throughitem:9:7and have readers pick a random suffix, spreading reads over eight shards. Writes must update or delete all copies. - Read from replicas for that shard, accepting slightly stale reads.
- Detection. Sample key access counts on clients or proxies so you find hot keys before they cause an outage.
The thundering herd, also called a cache stampede, happens when a popular key expires or is deleted and many requests miss at once. All of them hit the database for the same row. Fixes:
- Single flight or request coalescing. Within one process, concurrent misses for the same key share one database call. Across processes, use a short-lived lock or lease so one caller rebuilds and others wait briefly or read a stale copy.
- Stale-while-revalidate. Store a soft expiry inside the value. After it passes, serve the old value while one background request refreshes it.
- TTL jitter. Add a random spread to TTLs so keys written together do not all expire together.
- Probabilistic early refresh. As expiry approaches, each request has a growing chance of refreshing early, so one request usually rebuilds before the deadline.
These fixes overlap with what you would say about a rate limiter design, since both protect a fragile backend from bursts. If locking and coordination questions make you uneasy, review concurrency and multithreading interview questions first.
How Does Memcached vs Redis Fit Into Your Answer?
Memcached and Redis are both in-memory key-value stores, but they target different needs, and the interviewer may ask which one your design resembles.
| Dimension | Memcached | Redis / Valkey |
|---|---|---|
| Data model | Opaque strings | Strings, hashes, lists, sets, sorted sets, streams, more |
| Threading | Multithreaded | Command execution mainly on one thread, with optional I/O threads |
| Clustering | Client-side (consistent hashing in the client) | Built-in Redis Cluster with hash slots and failover |
| Replication | None built in | Primary-replica, asynchronous |
| Persistence | None | Optional RDB snapshots and AOF logs |
| Memory management | Slab allocator with segmented LRU | Configurable eviction policies, including LFU |
Licensing is a fair one-line aside. Redis changed its license in 2024, and the community responded with Valkey, a BSD-licensed fork hosted by the Linux Foundation. Check current license terms before recommending either in a real project.
What Follow-Up Questions Should You Expect?
These follow-ups come up often:
- How do you add a node without a miss storm? Migrate slots or ring arcs gradually, and let the new node warm up by reading from the old owner on misses before it takes full ownership.
- How do you run the cache across regions? Usually one cache per region with invalidations sent across regions from the database change stream. Cross-region writes to a shared cache add too much latency.
- How do you handle large values? Compress them, split them into chunks, or store a pointer to blob storage. Large values hurt tail latency for every other request on that node.
- What metrics would you alert on? Hit ratio, eviction rate, p99 latency, memory fragmentation, replication lag, and per-key request rate for hot key detection. This overlaps with the operational focus of a site reliability engineer interview.
- What if the whole cache cluster goes down? The database must survive a cold start. Use admission control, a rate limit on cache-miss traffic, and pre-warming from a key list.
- Can the cache be your source of truth? Only with persistence, synchronous replication, and
noeviction, at which point you are designing a database, not a cache.
For a full framework that applies to every prompt, read how to ace the system design interview, and use how to practice system design interviews to rehearse this answer out loud with a timer. Senior and staff candidates should expect deeper probing on failover and multi-region tradeoffs, which the staff engineer interview guide covers in more depth.
A strong 45-minute answer follows the section order above. If you can draw the ring, fill in the two tables from memory, and explain one failure scenario per pattern, you are ahead of most candidates.
Want a second brain in the room when the interviewer asks how your cache survives a primary failure mid-partition? TechScreen stays invisible during screen shares and gives real-time hints for system design and coding rounds. Try it with 3 free tokens.
Frequently Asked Questions
What is a distributed cache?
A distributed cache is an in-memory key-value store spread across many machines so that its total capacity and throughput exceed what one server can hold. Clients or a proxy route each key to the node that owns it, usually with consistent hashing or a fixed slot map. It sits in front of a slower system of record, such as a relational database, to cut read latency and database load. Redis, Valkey, and Memcached are the common building blocks.
How does consistent hashing help a distributed cache?
Consistent hashing places nodes and keys on the same hash ring, and each key belongs to the first node clockwise from its position. When a node joins or leaves, only the keys in its arc move, roughly 1/N of the data, instead of nearly every key as with hash(key) mod N. Virtual nodes give each server many small arcs, which smooths load and spreads a failed node's keys across all survivors rather than dumping them on one neighbor.
Which cache eviction policy should I choose in an interview?
Default to LRU or approximate LRU for general workloads because it is cheap and handles recency well. Choose LFU when a stable set of popular keys must survive bursts of one-off scans. Use TTL-based eviction when freshness matters more than popularity, and say noeviction only when the cache is really a primary store that must never silently drop data. The strong answer names the workload first, then the policy, then the failure it prevents.
What is the difference between cache-aside, write-through, and write-back?
In cache-aside, the application reads the cache, falls back to the database on a miss, and fills the cache itself. In write-through, every write goes to the cache and the database synchronously, so the cache stays fresh at the cost of write latency. In write-back, writes land in the cache and are flushed to the database later, which is fast but risks losing acknowledged writes if a cache node dies before flushing.
How do you prevent a thundering herd when a hot key expires?
Let only one request rebuild the value while others wait or serve stale data. Common tools are a per-key lock or lease, request coalescing inside each app server, serving a stale copy while a background job refreshes it, and adding random jitter to TTLs so many keys do not expire at the same moment. Facebook described a lease mechanism for exactly this problem in its Scaling Memcache paper.
Should I pick Redis or Memcached for a system design answer?
Pick Memcached when you need a simple, multithreaded cache of opaque strings and nothing else. Pick Redis or its fork Valkey when you need data structures like sorted sets or hashes, replication with automatic failover, persistence options, or built-in clustering. In most interviews the better move is to design the cache yourself and mention which product maps to your design, rather than answering with a product name.
How long should a distributed cache system design answer take?
Plan for about 45 minutes. Spend five on requirements and the API, five on rough capacity numbers, ten on partitioning and the high-level diagram, and the remaining time on replication, eviction, caching patterns, and hot keys. Leave a few minutes for follow-up questions. Interviewers care more about clear tradeoffs in the deep dives than about drawing every box.
Ready to use AI assistance in your next interview?
TechScreen is the invisible AI assistant trusted by engineers interviewing at Google, Meta, Amazon, and hundreds of other companies. Start with 3 free tokens — no credit card required.
Ace your next interview →