Design a Distributed Cache (Redis / Memcached Architecture)
Build a distributed in-memory cache: Consistent Hashing with virtual nodes, O(1) LRU eviction (HashMap + Doubly Linked List), Master-Replica replication, and gossip cluster state.
Distributed In-Memory Cache with Consistent Hash Ring & $O(1)$ LRU โก
Ketama consistent hashing with virtual nodes routing client keys to memory shards; each shard runs an O(1) LRU eviction engine.
01.1. Functional & Non-Functional Requirements
A Distributed In-Memory Cache stores frequently accessed key-value data in DRAM across a cluster of servers to achieve sub-millisecond read/write latencies and offload heavy database read IOPS.
Functional Requirements
- Core Operations: Support
GET(key),SET(key, value, ttl), andDEL(key)with sub-millisecond execution. - Eviction Policies: Support configurable eviction strategies when memory is full: LRU (Least Recently Used), LFU (Least Frequently Used), and FIFO.
- Configurable TTL: Automatic expiration of keys after a defined Time-To-Live.
- Data Structures: Basic string blobs, binary data, and optional collections (lists, sets, hashes).
Non-Functional Requirements
- Sub-Millisecond Latency: p99 read and write latency
< 1 ms. - Horizontal Elastic Scalability: Add or remove cache servers dynamically with minimal key remapping (
< 1/Nkeys remapped). - High Availability & Durability:
99.99\%uptime with automated master-replica failover. - Cache Miss Protection: Prevent stampedes, penetration, and avalanche failures.
02.2. Single-Node Storage Engine: $O(1)$ LRU Cache
To achieve strictly constant O(1) time complexity for both key lookups and memory evictions on a single server, combine two data structures in RAM:
- Hash Map (
std::unordered_map/ JavaHashMap):- Maps
key โ Pointer to Node in Doubly Linked List. - Provides instant
O(1)key lookup.
- Maps
- Doubly Linked List:
- Stores the cached key-value payload and maintains access recency order.
- Head: Most recently accessed item.
- Tail: Least recently used item (first candidate for eviction).
- Node removal and insertion at Head execute in
O(1)pointer adjustments without shifting arrays.
typescriptclass Node { key: string; val: any; prev: Node | null = null; next: Node | null = null; constructor(key: string, val: any) { this.key = key; this.val = val; } } class LRUCache { private capacity: number; private map: Map<string, Node> = new Map(); private head: Node = new Node("", null); private tail: Node = new Node("", null); constructor(capacity: number) { this.capacity = capacity; this.head.next = this.tail; this.tail.prev = this.head; } get(key: string): any { const node = this.map.get(key); if (!node) return null; this.moveToHead(node); return node.val; } set(key: string, val: any): void { if (this.map.has(key)) { const node = this.map.get(key)!; node.val = val; this.moveToHead(node); } else { if (this.map.size >= this.capacity) { const lru = this.tail.prev!; this.removeNode(lru); this.map.delete(lru.key); } const newNode = new Node(key, val); this.addNode(newNode); this.map.set(key, newNode); } } private addNode(node: Node) { node.prev = this.head; node.next = this.head.next; this.head.next!.prev = node; this.head.next = node; } private removeNode(node: Node) { node.prev!.next = node.next; node.next!.prev = node.prev; } private moveToHead(node: Node) { this.removeNode(node); this.addNode(node); } }
03.3. Cluster Partitioning: Consistent Hashing with Virtual Nodes
A single server cannot hold hundreds of gigabytes of cache data. Naive modulo hashing (hash(key) % N) causes catastrophic total cache invalidation when adding or removing a node, because almost 100\% of keys remap to different servers, triggering a database-crushing stampede.
The Consistent Hashing Ring Solution:
- Map both physical servers and cache keys onto a continuous
360^\circhash ring (0to2^{32}-1). - A key is routed to the first server encountered clockwise on the ring.
- Node Addition/Removal Impact: When a node is added or removed, only
1/Nfraction of keys are remapped.
Virtual Nodes (vNodes):
- Placing only 3 physical nodes on a ring causes non-uniform distribution (hot spots).
- Assign 100 to 200 Virtual Nodes (vNodes) per physical machine scattered uniformly across the ring.
- Result: Statistical standard deviation of key distribution drops to
< 3\%, ensuring balanced memory utilization across all cluster nodes.
04.4. Concurrency Architecture: Single-Threaded vs Multi-Threaded Engines
Understanding the concurrency architecture of the caching engine is critical for sizing and CPU core allocation:
1. Redis Model (Single-Threaded Event Loop)
- Uses non-blocking I/O multiplexing (epoll / kqueue / io_uring) on a single core for command execution.
- Why it is fast: Eliminates thread context switching, race conditions, and mutex lock contention. RAM access is so fast (
100 ns) that CPU is rarely the bottleneck. - I/O Threads (Redis 6.0+): Offloads network socket reading and writing to background worker threads while keeping core command execution strictly single-threaded.
2. Memcached Model (Multi-Threaded with Mutex Locks)
- Spawns multiple worker threads sharing a single memory address space.
- Employs fine-grained thread locking (LRU locks and Hash Table bucket locks).
- Utilizes multiple CPU cores effectively on large multi-socket server machines.
05.5. Cache Failure Modes & Production Mitigations
In production distributed caches, three distinct failure patterns can take down backend databases:
1. Cache Stampede / Thundering Herd
- Problem: A high-traffic key (e.g., home page product catalog with
50,000 QPS) expires. Thousands of simultaneous requests experience a cache miss at the exact same millisecond and all hammer the database at once. - Mitigation:
- Probabilistic Early Expiration (XFetch Algorithm): Recompute the cache in the background before it expires if
\Delta - \beta ร \ln(rand()) > TTL. - Mutex Lock (Single-Flight): The first cache-miss thread acquires a distributed lock to query the DB; all other threads wait or return stale cache data.
- Probabilistic Early Expiration (XFetch Algorithm): Recompute the cache in the background before it expires if
2. Cache Penetration
- Problem: Attackers query keys that do not exist in either the cache or the database (e.g.,
GET /users/invalid_id_-999). Every request bypasses the cache and hits the database. - Mitigation:
- In-Memory Bloom Filter: Place a Bloom filter in front of the cache to drop invalid IDs immediately.
- Null Value Caching: Cache empty results (
"user:-999" -> NULL) with a short TTL (e.g., 60 seconds).
3. Cache Avalanche
- Problem: Millions of keys are written with the identical TTL (e.g.,
1 hour) and all expire simultaneously, flooding the database. - Mitigation: Add random jitter to TTLs:
TTL = base_ttl + random_between(0, 300)seconds.
06.6. High Availability & Replication Topology
Master-Replica Pairs with Sentinel / Raft
- Each cache partition consists of a Primary Master and 1-2 Read Replicas.
- Writes hit the Primary; reads can be routed to replicas for increased read throughput.
- Replication is asynchronous for speed, with optional semi-synchronous replication for high durability.
- If the Primary crashes, a consensus quorum (Redis Sentinel or Raft-based Redis Cluster) automatically promotes a replica to Primary in
< 3 seconds.
โ๏ธArchitectural Trade-offs & Production Realities
Architectural Advantages
- Sub-millisecond read/write latency ($< 1\text{ms}$) by serving data directly from DRAM
- Consistent hashing with vNodes enables seamless horizontal cluster scaling with minimal key movement ($1/N$ remapping)
- O(1) LRU eviction maintains bounded memory usage without CPU spikes
Trade-offs & Constraints
- RAM is expensive compared to SSD/NVMe persistent storage ($~10\times$ cost differential)
- Asynchronous master-replica replication allows transient data loss on ungraceful primary crashes
AWS ElastiCache and Redis Cluster serve hundreds of millions of requests per second for enterprise platforms, using 16,384 hash slots, master-replica failover pairs, and client-side Ketama consistent hashing rings.
๐ฏ Staff+ Engineering Takeaways
- Consistent hashing with virtual nodes distributes keys evenly and prevents cluster-wide remapping during scaling.
- Single-node LRU combines a Hash Map and a Doubly Linked List for $O(1)$ lookups and evictions.
- Mitigate Thundering Herd using single-flight mutexes or probabilistic early expiration (XFetch).
- Use TTL jitter to prevent cache avalanches and Bloom filters to block cache penetration.
Topic Knowledge Assessment ๐ง
Step through 2 scenario questions to test your staff-level grasp.
Why are Virtual Nodes (vNodes) essential in a Consistent Hashing ring for a distributed cache?
How clear and staff-actionable was this system breakdown?