Design a Distributed Cache

⭐ Interview Importance: HIGH
⏱️ Revision Time: 4 min

Concept

The Problem: Design a distributed in-memory cache system like Redis or Memcached.

This question tests your knowledge of In-Memory Data Structures, Eviction Policies, and Consistent Hashing.

1. Requirements

Functional:

  • put(key, value)
  • get(key)
  • Fast eviction of old data when full.

Non-Functional:

  • Ultra-Low Latency (Sub-millisecond).
  • High Concurrency.
  • Distributed across multiple servers.

2. The Core Data Structure (LRU Cache)

A single cache server is fundamentally an LRU (Least Recently Used) Cache. It requires two data structures working together:

  1. Hash Map: Provides O(1)O(1) lookup speed for get(key).
  2. Doubly Linked List: Provides O(1)O(1) eviction speed. The most recently accessed items are at the “Head” of the list. The oldest items are pushed to the “Tail”.

The Workflow:
When a user calls get("user123"), the Hash Map finds the node in memory. Crucially, the system instantly snips that node out of its current position in the Doubly Linked List and moves it to the Head (marking it as the most recently used).
When RAM fills up, the system simply deletes the node currently sitting at the Tail of the list (O(1)O(1) time) to free up space.

3. High-Level Architecture (Distributed)

One server can hold perhaps 100GB of RAM. If we need to cache 5 Terabytes of data, we need 50 servers.

Consistent Hashing:
The Cache Router (or the Client SDK itself) uses Consistent Hashing to figure out which server holds the data. If Node B crashes, only the keys assigned to Node B are lost (Cache Misses). The requests are seamlessly rerouted to Node C, which will pull the data from the slow PostgreSQL database and repopulate the cache.

4. Availability vs Consistency (The Memcached vs Redis debate)

If Node B crashes, what happens?

  • The Memcached Approach (AP): We do not replicate data. If Node B crashes, the data is gone. The application will experience a sudden spike in “Cache Misses” and will fall back to querying the primary PostgreSQL database. This is cheap and fast, but puts the primary database at risk.
  • The Redis Approach (CP/Replication): Redis supports Master-Slave replication. Node B synchronously replicates its RAM to Node B-Replica. If Node B crashes, B-Replica instantly takes over, meaning 0 Cache Misses. This is safer but significantly more expensive (double the RAM costs).

Interview Questions

Q: A massive celebrity posts a photo. 10 million people request the photo’s metadata at the exact same millisecond. The metadata is NOT in the cache yet. What happens to the system, and how do you prevent it?
A: This is the Cache Stampede (Thundering Herd) problem.
Because the data is missing from the cache, all 10 million Node.js threads will instantly query the slow PostgreSQL database simultaneously to fetch it. The database will melt down and crash.
Fix: You must implement a Mutex Lock (or Promise Coalescing). The very first Node.js thread that encounters the Cache Miss acquires a distributed lock in Redis. The other 9,999,999 threads check for the lock, see it’s taken, and go to sleep for 50ms. The first thread queries PostgreSQL, puts the data into the Cache, and releases the lock. The sleeping threads wake up, check the Cache again, and successfully read the data without ever touching PostgreSQL.

Q: If you are building an LRU cache, why must you use a Doubly Linked List instead of a standard Array or Singly Linked List to track the order of items?
A: When an item is accessed via the Hash Map, it must be moved to the front (Head) of the list to mark it as “Recently Used”.

  • In an Array, deleting an item from the middle requires shifting every single subsequent element left (O(N)O(N)), which is too slow.
  • In a Singly Linked List, to delete a node, you must change the pointer of the previous node. But you only have a pointer to the current node. To find the previous node, you have to traverse the entire list from the Head (O(N)O(N)).
  • A Doubly Linked List gives you pointers to both the Next and Previous nodes, allowing you to snip the node out and reconnect its neighbors in perfect O(1)O(1) time.