Designing a Distributed Cache: The System Behind Your Session Store

You need sub-millisecond response times at scale. Your database cannot do this. You reach for a cache. Now design the rest.

In 2013, Facebook published a paper describing how Memcached served over a billion users. The cache tier handled over a trillion requests per day. Below it sat a MySQL cluster. Without the cache, that cluster would have collapsed under a fraction of the read load. The cache was not an optimisation. It was load-bearing infrastructure.

That is the frame for this article. A cache is not a transparent layer. It is an architectural component with its own failure modes and its own consistency model.


The problem you are actually solving

A high-read-QPS system needs data faster than the backing store can provide it. The backing store returns in tens of milliseconds. A cache returns in under one. The cache holds a working set of hot data in memory and serves reads without contacting the backing store.

The requirements are specific. Sub-millisecond read latency. Cache hit ratio above 90%. Graceful behaviour when the cache is cold. Graceful behaviour when a node fails. Predictable behaviour when the cache and the backing store disagree.

Every design decision below traces back to one of those four requirements.


Why the obvious approach fails

The naive approach is a dictionary in process memory. It works for a single-server application. It fails for a distributed system in two specific ways.

Each application server keeps an independent cache. A write invalidates the cache on the server that handled the write. It does not invalidate any other cache. Twenty servers hold twenty independent caches with different stale keys at different times. Consistency is undefined.

In-process caches also do not survive process restarts. Every deployment creates a cold cache. Every cold cache triggers a read storm on the backing database. That storm is FM7 — Thundering Herd, and it appears every time you push code.

The distributed cache exists to solve both problems.


The three topologies and when each applies

In-process caching embeds the cache inside the application. Fastest lookup. Independent caches. Loses everything on restart.

Sidecar caching runs a cache daemon on each host. The application talks to localhost. Slightly slower than in-process. Survives application restart. Still not shared across hosts.

The distributed cluster shares a cache across every application server. Multiple cache nodes hold a partition of the key space. The cluster is addressed using consistent hashing. A key always maps to the same node. Every server reads and writes the same cached value.

Multi-server systems need the cluster. The other two are optimisations on top.


Eviction: because memory is finite

Cache capacity is bounded. When the cache is full, something must be evicted.

LRU evicts the key accessed least recently. The assumption is that recently accessed data is likely to be accessed again. It fits most access patterns.

LFU evicts the key accessed the fewest times. It handles skewed access patterns better than LRU. It is more expensive to implement correctly.

TTL evicts every key at a fixed age regardless of access. It is simpler than LRU or LFU. It is the right choice for data with a known freshness bound — session tokens, rate limit counters, short-lived derived data.

Pick the policy that matches the access pattern. The wrong policy silently shrinks the hit ratio.


Write strategy: where the consistency choice actually lives

The write strategy determines when the cache is updated relative to the database. This is the decision that sets your consistency model.

Cache-aside is lazy. The application checks the cache. On a miss, it reads the database and writes the result back into the cache. The cache holds only data that has been read at least once. It is simple. It is prone to stampedes on cold start.

Write-through is synchronous. Every write updates both the cache and the database. The cache and the database are always consistent. Write latency doubles. The cache fills with data that may never be read.

Write-behind is asynchronous. Writes go to the cache immediately. The cache flushes to the database in the background. Write latency is the lowest of the three. If the cache node fails before flushing, the writes are lost.

This is AT4 — Precomputation vs On-Demand made concrete. Write-through precomputes. Cache-aside computes on demand. The right choice depends on the read-write ratio and whether reads can tolerate first-miss latency.


Replication: because a cache without it is a SPOF

Cache nodes are replicated in primary-replica pairs. The primary handles writes. The replica handles reads and takes over as primary if the primary fails. Redis Sentinel automates the failover. Redis Cluster provides sharding and replication together — each shard has a primary and one or more replicas.

Skip replication and you have built FM1 — Single Point of Failure. A cache node dies. Every key it held is suddenly missing. Every client for those keys misses simultaneously. Every miss hits the database. A node failure has become a thundering herd. A thundering herd has become a database outage.

Replication ensures a replica can immediately serve the missing keys. Nothing about this is optional at scale.


The failure that defines the system: thundering herd

A popular key expires. Every client waiting for it simultaneously misses. All of them query the database at once. The database was sized on the assumption that the cache would absorb this traffic. It cannot. Requests slow. Some time out. Some fail.

The cache repopulates from the first successful query. But the minutes it took to recover were minutes of degradation for every user.

Three mitigations exist. Lock-based population lets only one client fetch from the database while the others wait. Probabilistic expiry randomly extends the TTL near expiry so simultaneous expirations do not happen. Background refresh has a task rewrite hot keys before they expire, so they never go cold.

The probabilistic mitigation is worth stating directly. Measure the fetch time from the database. Pick a tunable beta. Compute an early-expiry threshold as current_time > (expiry - beta × fetch_time × log(random())). When that fires, one client refreshes early. The stampede never forms.

FM7 is not a failure of the cache. It is a failure of coordination between the cache and the database.


The stale-data failure most teams accept without knowing

Cache-aside with long TTLs causes FM4 — Data Consistency Failure. A user updates their profile. Every server serves the old profile for the next TTL seconds. Whether this is acceptable is an application question.

For a profile picture, a five-minute window is fine. For an inventory count during a flash sale, five seconds is a business incident. If the application cannot tolerate the staleness, the TTL must be short or the write path must invalidate the cached key explicitly.

The tradeoff has a name. AT1 — Consistency vs Availability. Write-through buys stronger consistency with higher write latency. Cache-aside with TTL accepts bounded staleness. Write-behind accepts potential loss for the lowest write latency. You are picking one.


How this evolves at scale

At 10x the initial size, the cluster grows to dozens of nodes. Cluster management — detecting failed nodes, redistributing keys, promoting replicas — must be automated. Redis Cluster or a managed service like ElastiCache does this.

At 100x, the working set may exceed a single cluster. Hierarchical caching helps. An L1 in-process cache sits in front of an L2 distributed cache. The hottest keys avoid the network round-trip entirely. CDN-level caching handles globally distributed read load for public content.

The hit ratio is the metric that tells you whether the design is working. Below 80% signals either a working set too large for available memory or an access pattern that caching cannot help — uniform distribution, no hot keys.


The signal that tells you this applies to your system

Your read QPS exceeds what the backing store can sustain at acceptable latency. Your hot keys are accessed far more often than your cold keys. Both conditions must hold. If your access is uniform, caching cannot help you — every key is equally cold.

If both conditions hold, the design questions are already waiting. Which topology? Which eviction policy? Which write strategy? What is your stampede defence? What is your replication factor? What happens when the cache is cold on Monday morning after a Sunday deploy?

The cache is not the answer to a slow database. The cache is a distributed system in front of the database, with its own answers required for each of those questions.

The full framework treatment — compression blocks, three-level exercises, and the complete AT/FM mapping — is in Book 4: System Design in Practice, Chapter 4 (Distributed Cache). Free chapter available at computingseries.com/books/book4.