TOPIC #234Advanced 10 min read

Design a Distributed Cache (Redis / Memcached Architecture)

CSD
CompleteSystemDesign Editorial
Report an issue
Key takeawayCore Architecture Summary

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.

Key Glossary Concepts in this TopicAll Glossary Terms
Interactive Lab · 🎯 Consistent Hash Cache RingFull lab guide

Redis-style Cluster: Consistent Hash Ring + LRU Economics

Hash 8,000 real keys onto the ring, add a node, and count exactly how many keys move compared with modulo hashing.

Hot working set: 40 GB (66.6M hot keys × 600 B)
Keys per node (ideal 2000)701 / 3243 / 3592 / 464
Load spread (max−min)3128 keysvNodes smooth the ring
Keys moved on scale-out2817 (35.2%)modulo hash would move 80%
Est. hit rate80%40/40 GB hot set resident
DB QPS when hot key expires1
PRODUCTION FAILURE MODES:

› Avalanche: synchronized TTLs expire together — jitter every TTL ±10–20% so misses spread over time.

› Penetration: attackers query non-existent keys, all missing to the DB — a Bloom filter of valid keys short-circuits them.

› Stampede: one 50,000-QPS key expires this millisecond — single-flight lets 1 request rebuild while the rest wait on it.

A single node serves keys via O(1) LRU: a HashMap for lookup plus a Doubly Linked List for recency, so get/put/evict never scan. The cluster layer then routes with a consistent-hash ring (client-side Ketama or 16,384 hash slots); each physical host occupies 50 ring positions, which is what keeps per-node key counts near the ideal 2000 instead of wildly uneven slices.

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.

Distributed In-Memory Cache with Consistent Hash Ring & $O(1)$ LRU ⚡
100%
Touchpad: Pinch to zoom • Drag to pan
Rendering visual architecture flowchart...

01.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

  1. Core Operations: Support GET(key), SET(key, value, ttl), and DEL(key) with sub-millisecond execution.
  2. Eviction Policies: Support configurable eviction strategies when memory is full: LRU (Least Recently Used), LFU (Least Frequently Used), and FIFO.
  3. Configurable TTL: Automatic expiration of keys after a defined Time-To-Live.
  4. 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/N keys remapped).
  • High Availability & Durability: 99.99\% uptime with automated master-replica failover.
  • Cache Miss Protection: Prevent stampedes, penetration, and avalanche failures.

02.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:

  1. Hash Map (std::unordered_map / Java HashMap):
    • Maps key ➔ Pointer to Node in Doubly Linked List.
    • Provides instant O(1) key lookup.
  2. 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.
typescript
class 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.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^\circ hash ring (0 to 2^{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/N fraction 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.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.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.

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.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
Production Implementation in Big Tech
Redis Cluster & Memcached (AWS ElastiCache)• Global In-Memory Data Tier

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 Check

Exercise 1 of 2 • Test your architectural comprehension.

Exercise 1 of 20 answered
1

Why are Virtual Nodes (vNodes) essential in a Consistent Hashing ring for a distributed cache?

Rate This Architecture ChapterFeedback & Rating

How clear and actionable was this distributed systems breakdown?

Interactive Engineering Workbenches: