Skip to main content

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

  1. Core Key-Value Operations: get(key), put(key, value, ttl), delete(key) with sub-millisecond latency.
  2. Configurable Eviction Policies: Evict cold items when memory limit is reached (LRU, LFU, W-TinyLFU, FIFO).
  3. Data Expiration (TTL): Keys expire automatically after their Time-To-Live (TTL) has elapsed.
  4. Horizontal Scalability: Add or remove cache nodes dynamically with minimal cache disruption (Consistent Hashing).
  5. High Availability & Replication: Primary-replica clustering with automated failover.

Non-Functional Requirements

  • Sub-Millisecond Read/Write Latency: P99 latency < 1ms for 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): 64Β GBΓ—0.80β‰ˆ51.2Β GBΒ usableΒ RAM/node64\text{ GB} \times 0.80 \approx 51.2\text{ GB usable RAM/node}.
    • Total Nodes = 10Β TB/51.2Β GBβ‰ˆ10\text{ TB} / 51.2\text{ GB} \approx 200 Cache Shards.
    • With 1 replica per shard (2x replication) β€…β€ŠβŸΉβ€…β€Š\implies 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
Client / AppBrowser / MobileAPI GatewayEnvoy / NGINXURL ServiceStateless Golang/JavaRedis CacheCluster (LRU)Primary DBPostgres / DynamoDBKafka / FlinkClick Analytics
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​

  1. The application client driver holds a local view of the Consistent Hashing Ring containing virtual nodes (tokens) for all primary cache instances.
  2. Client calls get("user:101").
  3. The client computes: hash=MurmurHash3("user:101")\text{hash} = \text{MurmurHash3}(\text{"user:101"}).
  4. Finds the first virtual node clockwise on the ring β€…β€ŠβŸΉβ€…β€Š\implies Node 4.
  5. 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​

  1. Node receives put(key, value, ttl).
  2. Checks current memory usage against max allocation (maxmemory).
  3. If memory is full:
    • Triggers the Eviction Engine (LRU / W-TinyLFU): Evicts victim keys from RAM.
  4. Stores key, pointer, metadata, and expiration timestamp in the in-memory hash table.
  5. 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 N=10N=10 nodes and 1 node crashes (N=9N=9):
    • 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 232βˆ’12^{32}-1 integer ring.
    • When a server is added or removed, only K/NK/N 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 β€…β€ŠβŸΉβ€…β€Š\implies OutOfMemoryError crash!
  • 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 jemalloc with automatic background defragmentation (activedefrag yes).
    • The engine copies allocated values to contiguous memory in the background while updating pointers without blocking reads.

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:​

  1. 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.
  2. Probabilistic Early Expiration (XFetch Algorithm):
    • Recompute the cache value before it actually expires using probabilistic math: Δ×β×ln⁑(rand())>TTLβˆ’now()\Delta \times \beta \times \ln(\text{rand}()) > \text{TTL} - \text{now}().
    • A background thread silently refreshes the key before the hard expiration boundary arrives.

5. Architectural Trade-Off Matrix

Design AreaOption AOption BSelected Choice & Rationale
Routing TopologyProxy-Based Routing (Twemproxy/Envoy)Client-Side Consistent Hash DriverClient-Side Driver: Direct socket communication eliminates a proxy network hop, shaving 0.5ms off every request.
Eviction AlgorithmStandard LRUW-TinyLFU (Window TinyLFU)W-TinyLFU: Solves the scan resistance flaw of standard LRU. Prevents full-cache evictions caused by one-off analytical scans.
Replication SemanticsSynchronous ReplicationAsynchronous Primary-ReplicaAsynchronous: 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.
πŸ“–
Track Page Progress0 / 635 Read
Knowledge Base Completion0%