Limited Offer

30% OFF Lifetime Access ($139) with code SYSTEM30

TOPIC #234Advanced 10 min read

Design a Distributed Cache (Redis / Memcached Architecture)

๐Ÿ’ก
Core 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

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%
Rendering visual architecture flowchart...

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

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

  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.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^\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.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.

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
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 Assessment ๐Ÿง 

Step through 2 scenario questions to test your staff-level grasp.

Question 1 of 20 answered
#1

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

Rate This Architecture Chapter4.9 / 5.0 (38 ratings)

How clear and staff-actionable was this system breakdown?