Designing a Distributed Cache.

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:
- Support
get,set, anddeleteoperations on key-value pairs. - Allow an optional expiration or time to live (TTL) for each entry.
- Evict entries using a Least Recently Used (LRU) policy when memory is full.
- 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:
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.
GET /cache/:keyA hit returns 200 with {"value": "..."}. A missing or expired entry returns 404.
DELETE /cache/:keyDeletion 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:
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.

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

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:
- We create a
Nodeobject containing thekey,value, andexpires_at - We add an entry to our hash table mapping the
keyto thisNode - We insert the
Nodeat the front of our doubly-linked list (right after the dummy head) - 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 :

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.

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.

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:
100,000 / 20,000 = 5 nodesNow 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:
ceil(1,024 / 28) = 37 nodes for one copySince 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.

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:
- First, the system monitors access patterns to detect hot keys.
- 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 copyuser:123#2-> Node B stores a copyuser:123#3-> Node C stores a copy
- These copies get distributed to different nodes via consistent hashing
- For reads, clients randomly choose one of the suffixed keys, spreading read load across multiple nodes.
- For writes, the system updates the copies asynchronously to achieve eventual consistency.

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:video123might be split into 10 shards:views:video123:1throughviews: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.

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:
- A hash table for storing our key-value pairs
- 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:
- Asynchronous replication for high availability and handling hot key reads.
- Consistent hashing for sharding and routing.
- Random suffixes for distributing hot key writes across nodes.
- Write batching and connection pooling for decreasing network latency/overhead.
Regardless of choices, the final design might look something like this:

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 ✌️