System DesignDistributed SystemsCachingLRUConsistent HashingReplication

Designing a Distributed Cache.

A client with a routing map contacts one of three shard primaries; each primary replicates asynchronously to two zone-separated replicas.
A distributed cache with sharded primaries and asynchronous replicas across availability zones.

A cache is simply a temporary storage space for fast reads and writes. The interesting part begins when the data no longer fits on one machine, a node fails, or one popular key gets more traffic than the rest of the cluster combined.

In this article, we’ll go through on how design a distributed cache system that can handle real production load and spans multiple nodes. Let’s get started.

What is a distributed cache?

A distributed cache stores key-value pairs in memory across multiple machines.
It is not limited by the resources of one machine and scales horizontally across multiple nodes to handle massive workloads.

The cache cluster works together to partition and replicate data, ensuring high availability and fault tolerance when individual nodes fail.

Requirements and scale

We need to:

  1. Support get, set, and delete operations on key-value pairs.
  2. Allow an optional expiration or time to live (TTL) for each entry.
  3. Evict entries using a Least Recently Used (LRU) policy when memory is full.
  4. Favor high availability, accepting temporarily stale reads.

The Scale we expect

Our expected scale is about 1 TiB of logical cache data and 100,000 requests per second at peak.

I'll use TiB here because the storage calculation below assumes 1,024 GiB.

The latency target needs a percentile and a workload to be meaningful; we'll evaluate it at p99 under the expected peak load.

Getting started with the design

I like to start with the core entities because they make it easier to reason about what the system stores and how requests change it.

For this design, the core entities are right there in front of our face! We're building a cache that stores key and values. That's it!.

Basically, we need to persist data that have keys, which are associated with values

API design

We have three operations. An HTTP interface could look like this:

Plain text
POST /cache/:key
Content-Type: application/json

{
  "value": "...",
  "ttl_seconds": 60
}

set replaces the value and its TTL. Omitting ttl_seconds means no expiration; if provided, it must be positive.

Plain text
GET /cache/:key

A hit returns 200 with {"value": "..."}. A missing or expired entry returns 404.

Plain text
DELETE /cache/:key

Deletion returns 204, including when the key is already absent. These are interface choices for our design; a compact cache protocol could expose the same operations with less serialization overhead.

A working cache on one node

Before introducing distributed systems, let's make the simplest version work:

Python
class Cache:
    def __init__(self):
        self.data = {}

    def get(self, key):
        return self.data.get(key)

    def set(self, key, value):
        self.data[key] = value

    def delete(self, key):
        self.data.pop(key, None)

At its core, a cache is just a hash table.
You could use Python's dict, Java's HashMap, and Go's map , etc. All of them give lookups and insert in O(1) time.

We can host this on one server and dispatch incoming requests to the corresponding method.

A client sends get, set, and delete requests to a cache node containing an in-memory hash table.
Figure 1. One node is enough to establish the basic request path: client → cache → hash table.

Adding expiration

To add expiration functionality we need to store an expiration timestamp alongside each value, then check it on reads. We’ll also need a way to clean up expired entries:

Python
def set(key, value, ttl_seconds=None):
    expires_at = (
        now() + ttl_seconds if ttl_seconds is not None else None
    )
    data[key] = Entry(value=value, expires_at=expires_at)

def get(key):
    entry = data.get(key)
    if entry is None:
        return MISS
    if entry.expires_at is not None and now() >= entry.expires_at:
        remove(key)
        return MISS
    return entry.value

To avoid the problem of our cache getting filled with expired entries, we will have a background process that periodically scans for and removes the expired entries using the cleanup function.

A client reads the cache while a separate background worker scans for expired entries.
Figure 2. Reads enforce expiration immediately; the worker reclaims memory from expired entries that receive no reads.

Adding LRU eviction

Now, we need to handle what happens when our cache gets full.

As stated earlier, we'll use the Least Recently Used (LRU) policy, which removes the entries that haven't been accessed for the longest time.

To build an LRU cache that allows us to do both, get and set operations in O(1) time, we combine two data structures:

  • A hash table for O(1) lookups.
  • A doubly linked list to track access order.

Let's see how it works:

  1. We create a Node object containing the key, value, and expires_at
  2. We add an entry to our hash table mapping the key to this Node
  3. We insert the Node at the front of our doubly-linked list (right after the dummy head)
  4. If we're at capacity, we remove the least recently used item from both the hash table and the linked list

The doubly-linked list maintains the exact order of access. The head of the list contains the most recently used items, and the tail contains the least recently used. When we need to evict an item, we simply remove the node right before our dummy tail.

Here’s a quick GIF that shows how LRU works for a cache of capacity = 3 :

Animated LRU example: reading A moves it to the front; inserting D then evicts B from a full three-entry cache.
Figure 3. Accessing A makes it most recent. Adding D to a full three-entry cache evicts B, the least recent entry. The hash table points to the same nodes held in the list.

The hash table gives us O(1) lookups, while the doubly-linked list gives us O(1) updates to our access order.

The best part about this implementation is that all operations (get, set, and even delete) remain O(1). When we access or add an item, we move it to the front of the list. When we need to evict, we remove from the back.

Great, now let's work on making our system work at scale.

High availability and fault tolerance

The key challenge to making a system highly available and fault tolerant is data replication. We need multiple copies of data spread across different regions. But it opens a new set of questions:

  • how many copies?
  • which nodes should store the copies?
  • how to sync everything?
  • what happens when nodes fail or can't communicate?

There are several well-established patterns for handling this. Each has its own trade-offs between consistency, availability, and complexity.

There are two approaches we can choose from:

Asynchronous replication

Here, the primary nodeapplies a write locally and acknowledges it without waiting for every replica. It also sends the update to its replicas through an ordered stream. Those replicas may briefly serve the previous value while they catch up.

The client writes to the primary, receives an acknowledgement, and the primary sends updates to two replicas asynchronously.
Figure 4. Replica delivery isn't on the acknowledgement path. Dashed arrows represent background replication; their completion may lag behind the client's response.

For our design requirements, this aligns well with eventual consistency.
It also offers higher availability, as writes can proceed even when replicas are down.

The main trade-off here is that readers can see stale values, and an acknowledged write can disappear if the primary fails before any replica receives it. Handling failovers introduces complexity and potential downtime.

Gossip Protocol (I’m serious).

In this approach, each node is equal and can accept both reads and writes. The nodes periodically exchange information about their state with randomly selected peers using gossip protocols.

The gossip protocol ensures that changes eventually reach all nodes in the cluster. This provides scalability and availability since there's no single point of failure. When a node receives a write, it can process it immediately and then asynchronously propagate the changes to its peers.

A client can write to any of three equal cache replicas. Bidirectional peer exchanges asynchronously propagate updates between replicas.
Figure 5. A client can write to any replica. Replicas exchange updates asynchronously with peers; the arrows show example peer exchanges.

Although this approach offers great scalability and availability, it comes with some significant challenges. The implementation is more complex than other approaches since each node needs to maintain connections with multiple peers and handle conflict resolution. Additionally, careful consideration must be given to conflict resolution strategies when concurrent updates occur at different nodes.

Amazon's Dynamo paper uses gossip for membership and failure detection, alongside separate mechanisms for replica synchronization.

Scaling the cluster

To handle the scale of managing a terabyte of storage and 100k requests per second, We need to spread the load across multiple machines to maintain low latency and high availability.

First, Let's estimate capacity from both throughput and memory.

Assume a benchmark shows that one host can sustain 20,000 requests per second for our request mix and latency target and gnoring replication overhead for the moment:

Plain text
100,000 / 20,000 = 5 nodes

Now assume a 32 GiB RAM instance offers 28 GiB of effective entry capacity, after reserving memory for the process, entry metadata, hash table, linked list, and replication buffers:

Plain text
ceil(1,024 / 28) = 37 nodes for one copy

Since memory dominates this estimate, we should provision based on our storage needs. This higher count will also give us room for handling throughput requirements.

How to distribute keys evenly across nodes?

Consistent Hashing.

If you've designed systems before, you've probably read about consistent hashing.

It's a technique that helps us distribute keys across our cache nodes while minimizing the number of keys that need to be remapped when nodes are added or removed.

So, instead of using hash(key) % 5, we use a consistent hashing function like MurmurHash to get a position on the circle. Given that position, we move clockwise around the circle until we find the first node, this node should store the key-value pair we're looking for.

A key lies between C and A on a hash ring, and a clockwise arrow points to A as its owner.
Figure 6. The first node token clockwise from the key owns it. Adding a token changes ownership of the adjacent range rather than remapping the whole dataset.

What if there's a hot key?

Hot keys are a common challenge in designing systems. They occur when certain keys receive disproportionately high traffic compared to others. Eg: a viral tweet, ticket count in flash booking of a concert.

When too many requests concentrate on a single shard holding these popular keys, it creates a hotspot that can degrade performance for that entire shard.

There are two different problems: hot reads, where many clients ask for the same value, and hot writes, where many updates target the same logical item.

Scaling hot reads

For a read-heavy keys, we can create extra copies and spread reads across them. We can use healthy replicas first, then add selective copies if the existing replica set isn't enough.

Here's how it works:

  1. First, the system monitors access patterns to detect hot keys.
  2. When a key becomes "hot", instead of having just one copy as user:123, the system creates multiple copies with different suffixes:
    • user:123#1 -> Node A stores a copy
    • user:123#2 -> Node B stores a copy
    • user:123#3 -> Node C stores a copy
  3. These copies get distributed to different nodes via consistent hashing
  4. For reads, clients randomly choose one of the suffixed keys, spreading read load across multiple nodes.
  5. For writes, the system updates the copies asynchronously to achieve eventual consistency.
A client chooses one of three copies of the same value: user:123#1 on Node A, user:123#2 on Node B, or user:123#3 on Node C.
Figure 7. Each read chooses one suffixed key: user:123#1, user:123#2, or user:123#3. The copies hold the same value, spreading reads across nodes.

This approach is specifically designed for read-heavy hot keys. If you have a key that's hot for both reads and writes, this approach can actually make things worse due to the overhead of maintaining consistency across copies.

Scaling hot writes

Hot writes are a similar problem to hot reads, just that they're a little more complex and have a different set of trade-offs.

We can solve this in two ways:

  • Write batching. Combine compatible updates before sending them. Ten counter increments, for example, can become one increment by ten if the protocol supports an atomic increment. A longer batching window reduces request volume but delays visibility. A 50–100 ms window would already violate our ordinary sub-10-ms write target, so use it only for a separate workload that explicitly accepts that delay. Retries also need deduplication if applying a batch twice would change the result. This approach works best for metrics and counters where eventual consistency is acceptable, but may not be suitable for scenarios requiring immediate write visibility.
  • Sharding a hot key with suffixes. Split a hot key into multiple sub-keys distributed across different nodes. For eg, if we count views for videos, a hot counter key views:video123 might be split into 10 shards: views:video123:1 through views:video123:10.
    On write, send the update to increment to one sub-key selected randomly, then sum all ten on reads. At 10,000 writes per second, using 10 shards would reduce the per-shard write load to roughly 1,000 operations per second.

    We trade cheaper writes for more expensive reads, which may observe different moments in time. This works for aggregatable operations like counters and metrics, not arbitrary replacements or updates requiring strict global ordering.
A write randomly selects one counter shard, shown updating views:video123:2. A read sums partial counts from all ten shards, views:video123:1 through views:video123:10.
Figure 8. Each write increments one randomly selected shard; a read sums all ten partial counters. Shards 3–9 are omitted for clarity.


Keeping latency low

When you move from a single node to a distributed system with multiple nodes, performance isn’t just about how fast a hash table runs in memory. It’s about how quickly clients can find the right node, how efficiently multiple requests are bundled together, and how we avoid unnecessary network hops.

The good news is that we've already made useful choices.

With Consistent Hashing clients route directly from a local map, avoiding a routing lookup on every request. Although stale map may still require a refresh and retry.

Pipelining / Batching can send several independent requests without waiting for each response before sending the next. This reduces repeated round-trip waits and improves throughput.

One last thing is connection pooling. Constantly tearing down and re-establishing network connections wastes time. Instead of making a fresh connection for every request, clients should maintain a pool of open, persistent connections. This ensures there’s always a ready-to-use channel for requests, removing expensive round-trip handshakes and drastically reducing tail latencies (like those p95 and p99 response times that can make or break user experience).

Finally, we should always setup monitoring and observability to measure p99 latency while replicas catch up, nodes fail, and shards rebalance.

In production, everything that can go wrong, will go wrong. ~ Moyez.

Final design

Okay, tying it all together, on each node, we'll have two data structures:

  1. A hash table for storing our key-value pairs
  2. A linked list for our LRU eviction policy

When it comes to scaling, it depends on which of the above approaches you choose, as you can end up with a slightly different design.

My choice is:

  1. Asynchronous replication for high availability and handling hot key reads.
  2. Consistent hashing for sharding and routing.
  3. Random suffixes for distributing hot key writes across nodes.
  4. Write batching and connection pooling for decreasing network latency/overhead.

Regardless of choices, the final design might look something like this:

A client with a routing map contacts one of three shard primaries; each primary replicates asynchronously to two zone-separated replicas.
Figure 9. The routing map selects a shard's current primary for writes. Each row is one shard, with copies on separate nodes across three zones. The inset shows the structures inside every node; dashed arrows show asynchronous replication.

The diagram shows a few logical shards, not the full fleet. A physical node can host multiple shard roles, and clients can read replicas when stale reads are acceptable.

On a hit, the selected node checks expiration, updates local recency, and returns the value. On a write, the primary updates its local entry and replicates the change without waiting for all copies. On a miss, the application loads from its durable source and repopulates the cache.

The result is a cache that can grow beyond one machine and recover from individual node failures. Its guarantees follow from the workload: data is reconstructable, stale reads are acceptable, and the fast path stays local to a region. Change those requirements, and the replication, failover, and write strategy need to change with them.


Alright folks, this is it. Considering the advancements in AI, I will be posting fewer techincal articles now. Future articles will focus more on agentic coding and techniques to leverage these AI models to the fullest.

Follow me on X and LinkedIn for similar content.

Peace out ✌️