Design a Distributed In-Memory Cache Like Redis
A distributed in-memory cache (e.g., Redis Cluster, Memcached, Amazon ElastiCache) provides ultra-low latency, sub-millisecond key-value storage in RAM. Caching sits between application services and durable databases to absorb massive read spikes, reduce database CPU load, and accelerate API response times.
1. Understanding the Problem
Functional Requirements
- Core Key-Value Operations:
get(key),put(key, value, ttl),delete(key)with sub-millisecond latency. - Configurable Eviction Policies: Evict cold items when memory limit is reached (LRU, LFU, W-TinyLFU, FIFO).
- Data Expiration (TTL): Keys expire automatically after their Time-To-Live (TTL) has elapsed.
- Horizontal Scalability: Add or remove cache nodes dynamically with minimal cache disruption (Consistent Hashing).
- High Availability & Replication: Primary-replica clustering with automated failover.
Non-Functional Requirements
- Sub-Millisecond Read/Write Latency: P99 latency
< 1msfor all operations. - High Throughput: Support Millions of operations per second across the cluster.
- Memory Efficiency: Minimal memory overhead per key; zero memory fragmentation crashes.
- Fault Tolerance: Cluster continues serving traffic during node failures without data corruption or cascading crashes.
Capacity Estimations & Sizing
- Total Ingress Throughput: 5 Million operations/sec (80% reads, 20% writes).
- Total Cached Data: 10 Terabytes (TB) of active cached objects.
- Server Sizing:
- Cloud cache nodes equipped with 64 GB RAM and 10 Gbps network interfaces.
- Safe memory utilization cap (80% to avoid OS swapping): .
- Total Nodes = 200 Cache Shards.
- With 1 replica per shard (2x replication) 400 Total Nodes.
2. The Set Up
Eviction Policies & Data Structures
βββββββββββββββββββββββββββββ¬βββββββββββββββββββββββββββββββββββ¬βββββββββββββββββββββββββββββββββββ
β EVICTION POLICY β DATA STRUCTURE β COMPLEXITY & BEHAVIOR β
βββββββββββββββββββββββββββββΌβββββββββββββββββββββββββββββββββββΌβββββββββββββββββββββββββββββββββββ€
β 1. LRU (Least Recently β Hash Map + Doubly Linked List β O(1) get, O(1) put. β
β Used) β β Vulnerable to full-cache flushes β
β β β during one-off table scans. β
βββββββββββββββββββββββββββββΌβββββββββββββββββββββββββββββββββββΌβββββββββββββββββββββββββββββββββββ€
β 2. LFU (Least Frequently β Hash Map + Frequency Doubly β O(1) get, O(1) put. β
β Used) β Linked Lists β Historical bias: Old popular β
β β β items linger forever (decay req) β
βββββββββββββββββββββββββββββΌβββββββββββββββββββββββββββββββββββΌβββββββββββββββββββββββββββββββββββ€
β 3. W-TinyLFU (Window β TinyLFU (Count-Min Sketch) + β Near-optimal hit ratios! Used by β
β TinyLFU) β SLRU (Segmented LRU) β Caffeine Cache and modern Redis. β
βββββββββββββββββββββββββββββ΄βββββββββββββββββββββββββββββββββββ΄βββββββββββββββββββββββββββββββββββ
The Client Driver & Protocol API
// Standard Cache Operations
interface DistributedCacheClient {
byte[] get(String key);
boolean put(String key, byte[] value, long ttlSeconds);
boolean delete(String key);
}
3. High-Level Design
Distributed In-Memory Cache Cluster ArchitectureInteractive Topology
Read QPS
100K/sec
Write QPS
1K/sec
Latency SLA
< 15ms
5-Yr Storage
~15 TB
Active Scenario: User submits long URL -> Token Generator (KGS) allocates Base62 ID -> Writes to DB & Warm Cache
Interactive Component Inspector
Click any architecture node on the SVG canvas to view under-the-hood engine mechanics, failure gotchas, and runtime tags.
Core Subsystems & Operational Flow
1. Consistent Hashing Client Routingβ
- The application client driver holds a local view of the Consistent Hashing Ring containing virtual nodes (tokens) for all primary cache instances.
- Client calls
get("user:101"). - The client computes: .
- Finds the first virtual node clockwise on the ring Node 4.
- Issues a direct TCP socket call to Node 4 over a persistent connection pool in 0.4ms.
2. Storage & Eviction Inside a Single Cache Nodeβ
- Node receives
put(key, value, ttl). - Checks current memory usage against max allocation (
maxmemory). - If memory is full:
- Triggers the Eviction Engine (LRU / W-TinyLFU): Evicts victim keys from RAM.
- Stores key, pointer, metadata, and expiration timestamp in the in-memory hash table.
- Asynchronously streams the write command to connected replica nodes over replication sockets.
4. Potential Deep Dives & Bottlenecks
Deep Dive 1: Consistent Hashing with Virtual Nodes
Why does naive modulo hashing (hash(key) % N) fail in distributed caching?
- The Modulo Catastrophe:
- If we have nodes and 1 node crashes ():
- Almost 100% of keys remap to different nodes (
hash(key) % 9 != hash(key) % 10)! - Result: Immediate 100% cache miss storm across the entire company. Database receives 100x traffic surge and collapses instantly (Cascading Failure).
- Consistent Hashing Solution:
- Map keys and server nodes onto a circular integer ring.
- When a server is added or removed, only keys are remapped!
- Virtual Nodes (Preventing Non-Uniform Hotspots):
- Placing physical nodes directly on the ring causes uneven partition splits.
- By assigning 100 to 256 virtual nodes per physical server (e.g.
Server1#1,Server1#2), keys are distributed uniformly across all physical nodes within a 2% variance.
Consistent Hash Ring (0 to 2^32 - 1)
[Node 1#V1]
β
βββββββββββββββββ΄ββββββββββββββββ
β β
[Node 3#V2] [Node 2#V1]
β β
β Key "order:99" β
β hashes here βββββββββΊβ (Routed to Node 2)
β β
[Node 2#V2] [Node 1#V2]
β β
βββββββββββββββββ¬ββββββββββββββββ
β
[Node 3#V1]
Deep Dive 2: Memory Management & Fragmentation (Slab Allocation vs jemalloc)
Why does continuous allocation of dynamic key sizes crash cache servers?
- External Memory Fragmentation: Repeatedly allocating and freeing random-sized keys (100 bytes, 4 KB, 1 MB) creates holes in physical memory. The OS reports 10 GB of free RAM, but no contiguous block exists to allocate a 50 KB object
OutOfMemoryErrorcrash! - Slab Allocation (Memcached Model):
- Pre-allocates memory into fixed 1 MB pages.
- Divides pages into uniform chunks called Slab Classes (e.g. Class 1: 64 bytes, Class 2: 128 bytes, Class 3: 256 bytes).
- A 100-byte object is stored in Slab Class 2.
- Completely eliminates external memory fragmentation.
- jemalloc & Active Defragmentation (Redis Model):
- Uses
jemallocwith automatic background defragmentation (activedefrag yes). - The engine copies allocated values to contiguous memory in the background while updating pointers without blocking reads.
- Uses
Deep Dive 3: Cache Stampede (Thundering Herd) & Single-Flight Locking
What happens when a popular hot key ("breaking:news") expires under 50,000 requests per second?
50,000 Concurrent Requests arrive for Key "breaking:news"
β
βΌ
Key expired in Cache! (50,000 Cache Misses!)
β
βΌ
ALL 50,000 REQUESTS HIT THE PRIMARY DATABASE SIMULTANEOUSLY!
β DATABASE CONNECTION EXHAUSTION & CRASH!
Production Defenses:β
- Single-Flight / Mutex Locking (
Go singleflight/ Distributed Lock):- When a cache miss occurs, only one single worker acquires a local mutex to fetch the data from the database and populate the cache.
- The remaining 49,999 requests wait for the mutex or subscribe to the in-flight result, completely protecting the database.
- Probabilistic Early Expiration (XFetch Algorithm):
- Recompute the cache value before it actually expires using probabilistic math: .
- A background thread silently refreshes the key before the hard expiration boundary arrives.
5. Architectural Trade-Off Matrix
| Design Area | Option A | Option B | Selected Choice & Rationale |
|---|---|---|---|
| Routing Topology | Proxy-Based Routing (Twemproxy/Envoy) | Client-Side Consistent Hash Driver | Client-Side Driver: Direct socket communication eliminates a proxy network hop, shaving 0.5ms off every request. |
| Eviction Algorithm | Standard LRU | W-TinyLFU (Window TinyLFU) | W-TinyLFU: Solves the scan resistance flaw of standard LRU. Prevents full-cache evictions caused by one-off analytical scans. |
| Replication Semantics | Synchronous Replication | Asynchronous Primary-Replica | Asynchronous: Synchronous replication forces writes to wait for network round-trips to replicas, destroying sub-millisecond write SLAs. |
6. What is Expected at Each Level?
Mid-Level (L4 / IC4)
- Understands the need for caching in front of databases.
- Implements basic in-memory LRU cache using a Hash Map and Doubly Linked List.
- Explains basic TTL and cache expiration.
- Designs basic Cache-Aside read/write flows.
Senior (L5 / IC5)
- Details the Consistent Hashing algorithm and why virtual nodes are essential to prevent partition skew.
- Explains memory fragmentation and compares Slab Allocation (Memcached) vs jemalloc defragmentation (Redis).
- Solves Cache Stampede / Thundering Herd using single-flight mutexes or probabilistic early expiration (XFetch).
- Implements primary-replica failover mechanisms with automated health-check heartbeats.
Staff+ (L6 / Principal)
- Designs multi-datacenter cache synchronization: Handling cross-region cache replication without causing stale read loops or split-brain overwrites during network partitions.
- Details kernel networking optimization: Bypassing Linux kernel network stack overhead using kernel-bypass networking (DPDK / eBPF / XDP) to achieve 10M+ packets/sec per node.
- Evaluates persistent memory (PMEM / Intel Optane) trade-offs for instant cold-start cache warming after cluster restarts.
