Design a Distributed Cache
DifficultyHard | HelloInterview: problem breakdown
Problem Statement
Design a distributed in-memory cache (like Redis Cluster or Memcached) that stores key-value pairs, supports high throughput reads/writes with low latency, and scales horizontally.
๐Real-world: The unsung hero of this design is consistent hashing, which came out of a 1997 MIT paper and was originally built to solve exactly this โ distributing keys across caching nodes so that adding/removing a node only moves ~1/N of the keys instead of reshuffling everything (the same idea later powering Amazon's Dynamo and DynamoDB). The cautionary security tale for caches: for years tens of thousands of Redis and Memcached instances sat exposed to the internet with no authentication (Redis historically had no auth by default and trusted the network), leading to mass data theft and, worse, Memcached becoming the vehicle for record-breaking DDoS amplification in 2018 โ a ~1.3 Tbps attack on GitHub โ because an unauthenticated UDP cache reflects huge responses to spoofed victims. The lessons: bind to localhost/private networks, require auth, disable UDP, and never treat "it's just a cache" as "it doesn't need security." Also know the failure modes this design must handle: cache stampede, hot keys, and thundering herd.
Requirements
Functional
get(key) โ valueput(key, value, ttl)delete(key)- Configurable eviction policies (LRU, LFU, TTL)
- Cluster mode: data distributed across nodes
Non-Functional
- P99 latency: < 1ms for get/put
- Throughput: 1M ops/s per cluster
- High availability: no single point of failure
- Horizontal scaling: add nodes without downtime
- Consistency: eventual is acceptable (caches can be stale)
Core Design
Single Node Cache
Start simple: an in-memory hash map with LRU eviction.
HashMap<String, Node> map โ O(1) lookup
DoublyLinkedList list โ O(1) LRU move-to-front / evict-from-tail
get(key):
node = map.get(key)
if expired: evict, return null
move node to head of list
return node.value
put(key, value, ttl):
if key exists: update value, move to head
else:
if at capacity: evict tail
create node, add to head, add to mapConsistent Hashing for Distribution
Naive approach: hash(key) % N โ when N changes (node added/removed), almost all keys need to be remapped. Bad for a live system.
Consistent hashingplace nodes and keys on a conceptual ring (0 โ 2^32). A key belongs to the first node clockwise from its hash position.
Ring: โโโโโ Node A (90ยฐ) โโโโ Node B (180ยฐ) โโโโ Node C (270ยฐ) โโโโ Key K hashes to 120ยฐ โ assigned to Node B (next clockwise) Add Node D at 150ยฐ: Only keys between 120ยฐ-150ยฐ (previously B's) reassign to D All other keys: unchanged
Virtual nodeseach physical node occupies multiple positions on the ring (e.g., 150 virtual nodes per physical node). This ensures uniform key distribution even with heterogeneous hardware.
Replication for Availability
Each key is replicated to the next N nodes clockwise on the ring (N = replication factor, typically 3).
primary node receives write, replicates to N-1 replicas
can read from any replica (eventual consistency) OR require quorum (stronger consistency)
when a node fails, its virtual nodes are taken over by successors; replicas ensure no data loss
Architecture
Client SDK (consistent hash โ pick node)
โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ Cache Node A โ Cache Node B โ
โ [key1, key3] โ [key2, key4] โ
โ replicas of B โ replicas of C โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ Cache Node C โ Cache Node D โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
Coordinator / Config Server:
- Tracks live nodes
- Distributes ring membership to clients
- Health checksClient-side vs Proxy-side Routing
Client-side routing (Redis Cluster, Memcached): client library implements consistent hashing, connects directly to correct node. No proxy hop = lower latency. Client must know cluster topology.
Proxy-side routing (e.g., Twemproxy): all clients connect to a proxy; proxy routes to correct node. Simpler clients; proxy is a potential bottleneck.
Eviction Policies
| Policy | Description | Best for |
|---|---|---|
| LRU | Evict least recently accessed | General cache; temporal locality |
| LFU | Evict least frequently accessed | Stable hot data; celebrity keys |
| TTL-based | Evict expired keys proactively | Time-bounded data (sessions, tokens) |
| Random | Evict random key | Very simple; approximate LRU at low cost |
Active vs lazy expiration
- Lazy: check TTL on access, evict if expired
- Active: background thread periodically scans and evicts expired keys (prevents memory leak from never-accessed expired keys)
Key Design Decisions
| Decision | Choice | Reason |
|---|---|---|
| Key distribution | Consistent hashing with virtual nodes | Minimize remapping on membership change |
| Replication | 3 replicas, async | Availability without strong consistency overhead |
| Consistency | Eventual | Caches are acceptable to be slightly stale |
| Eviction | LRU + TTL | Most common access pattern |
| Routing | Client-side | Lower latency, no proxy bottleneck |
| Serialization | MessagePack or Protobuf | Compact binary; faster than JSON |
Security Considerations
| Threat | Mitigation |
|---|---|
| Cache poisoning | Validate data written to cache; cache should not be the source of truth for auth decisions |
| Cache key enumeration | Use opaque/hashed cache keys, not predictable patterns |
| Sensitive data in cache | Encrypt sensitive values (PII, tokens); set short TTL |
| Unauthorized cache access | Network-level access control (VPC, security groups); auth tokens for cache access |
| Cache side-channel | Cache timing attacks can leak whether a key exists; use consistent-time responses |
| Denial of service | Rate limit per client on cache access; limit value sizes; reject excessively large values |
| Cache stampede on invalidation | Probabilistic early expiration; mutex lock on miss; staggered TTLs |
AuthenticationRedis in production should always have requirepass or ACL-based auth, and should never be exposed on a public interface without authentication.
Interview Tips
walk through the data structure before scaling.
spend time on why % N is bad and how the ring fixes it.
don't skip this, it's why consistent hashing is practical.
connect this to CAP theorem: you're trading consistency for availability and partition tolerance.
"how do you handle cache invalidation at scale?" โ event-driven invalidation via message queue, vs TTL, vs version keys.