Distributed Caching Patterns & Eviction
The Product Page That Melted Redis
Test your architecture intuition: Pitch a 7-axis solution, survive two aggressive reviewer objections, and inspect the staff-level Teacher Gold Answer.
1. What It Is & Why It Exists
The Core Problem
Persistent storage engines (relational databases, document stores, and distributed filesystems) are fundamentally bounded by disk I/O, disk seek latencies, B-Tree index locking, and multi-tenant connection pools. A database query typically takes several milliseconds or more once parsing, locking and queueing are counted (measure your own).
In high-throughput distributed applications handling hundreds of thousands of read requests per second, querying disk-bound databases directly causes:
- Connection Pool Exhaustion: Available database worker threads and TCP sockets saturate within milliseconds.
- CPU & Lock Contention: Intensive query parsing, table lock contention, and buffer pool churn degrade database throughput.
- Cascading Latency Cliffs (Tail Latency Explosion): P99 and P99.9 latency spikes exponentially, causing upstream service timeouts, client retry storms, and total system brownouts.
The Breakdown: The Unshielded Database Collapse
Concrete Proof: Unshielded Database vs. Distributed Cache
Consider an e-commerce platform processing during a promotional event, where access follows an 80/20 Pareto distribution (80% of queries target 20% of catalog items):
| Metric / Dimension | Direct Database Architecture (No Cache) | Distributed In-Memory Cache Tier ( Hit Ratio) |
|---|---|---|
| Storage Medium | NVMe SSD / Disk Buffer Pool | Volatile RAM (Redis / Memcached) |
| Access Latency | Several milliseconds per query (assumed for this example) | Sub-millisecond: one network round trip plus microseconds of work |
| QPS Capacity per Node | (Database primary; assumed for this example, measure yours) | Often 100,000+ simple operations/s per node (In-memory node; measure, since it depends on node size, command mix and payload) |
| Database Load at 200k QPS | 🚨 Connection pool exhausted | 🛡️ ( traffic absorbed in RAM) |
| Database Fleet Sizing | Requires 40-100 read replicas at 2,000-5,000 QPS each | Single primary + 2 read replicas |
| P99 Read Latency | 🚨 Climbs without bound once the database saturates (queueing backlog) | ⚡ Close to the cache's own latency, as long as the 5% of misses fit the database |
The Database Throughput Wall: Relational databases cannot scale reads indefinitely via read replicas due to primary replication lag, cross-AZ synchronization overhead, and connection limits. Past that point, a distributed cache is usually the cheapest way to scale reads further.
The First-Principles Solution: In-Memory Caching
In-memory caching places high-speed, volatile RAM stores (e.g., Redis, Memcached, DynamoDB DAX) between the application tier and the primary persistence layer. Accessing RAM takes ; an SSD read takes ; a database query takes several milliseconds once parsing, locking and queueing are counted. A networked cache still pays a network round trip (a fraction of a millisecond), so the win over the database is roughly in latency and far more in capacity per node.
Synthesizing vector architecture diagram...
Both panels start with 100,000 clients reading through the same app fleet. In the "Broken Baseline" panel, every read becomes a database query, so the database gets all 100,000 queries per second, runs out of connections and CPU, and latency climbs as queries queue. In the "Production Standard" panel, the app checks the RAM cache first (step 1); 95% of reads are hits served in under a millisecond, and only the 5% misses go on to the database (step 2), which now carries a twentieth of the load and has headroom. Every point of hit rate matters: going from 95% to 90% doubles database load, which is why hit rate is the cache metric to watch.
2. Core Mechanics & Caching Design Patterns
Caching Design Patterns Deep-Dive
1. Cache-Aside (Lazy Loading)
- Read Path: The application queries the cache first. On a cache hit, data is returned immediately. On a cache miss, the application queries the database, writes the retrieved value to the cache with a Time-To-Live (TTL), and returns data to the client.
- Write Path: The application updates the database first, then evicts (deletes) the corresponding key from the cache (
DEL key). - Trade-Offs:
- Pros: Memory-efficient (only queried data is cached); resilient to cache node failures (cache misses gracefully fall back to the DB).
- Cons: Initial read experiences cache miss latency; vulnerable to dual-write race conditions if not paired with cache eviction.
2. Write-Through
- Mechanics: The application writes directly to the caching layer. The cache synchronously writes data to the underlying database before returning success to the client.
- Trade-Offs:
- Pros: With a single writer the cache stays in step with the database, and recently written keys are hits until evicted or expired. Two concurrent writers can still land in the cache in the opposite order from the database, so writes need a version check.
- Cons: High write latency (incurs RAM + DB round-trip latency); pollutes cache with data that may never be read again.
3. Write-Behind (Write-Back)
- Mechanics: The application writes directly to the cache, which immediately acknowledges success (). An asynchronous background worker batches dirty cache entries and writes them to the database periodically.
- Trade-Offs:
- Pros: Extreme write throughput; absorbs write bursts and collapses multiple writes to the same key into a single DB write.
- Cons: High risk of data loss—if the cache node crashes before flushing dirty pages, uncommitted updates are permanently lost.
4. Refresh-Ahead (Proactive Caching)
- Mechanics: Keys that are still being read are reloaded from the database shortly before their TTL expires (by a background job, a library's refresh-after-write setting, or probabilistic early refresh such as XFetch below), so readers of hot keys never see the miss.
- Trade-Offs:
- Pros: Eliminates cache miss latency for hot keys.
- Cons: Mispredicting access patterns causes unnecessary database load refreshing cold keys.
Synthesizing vector architecture diagram...
Each panel is one pattern; follow its numbered arrows. In the "Cache-Aside" panel, the app reads the cache (1), and on a miss queries the database itself (2) and writes the result into the cache with a TTL (3), so only data that is actually read gets cached. In the "Write-Through" panel, the app writes to the cache (1) and the cache synchronously writes to the database (2) before answering, so the cache follows every write (with one writer at a time) but every write pays both latencies. In the "Write-Behind" panel, the cache acknowledges the write immediately (1) and flushes batches to the database later (2), giving the fastest writes but losing anything not yet flushed if the cache node dies. Choose by what you can tolerate: stale reads (cache-aside), slow writes (write-through), or lost writes (write-behind).
Formal Mathematical Formulations
1. Effective Access Time (EAT)
The performance of a tiered memory architecture is governed by its Cache Hit Ratio ():
Where is in-memory retrieval latency (e.g. ) and is primary database query latency (e.g. ); both are example values.
The 90% vs 99% Rule:
- At : .
- At : ( lower average latency).
2. Probabilistic Early Expiration (The XFetch Algorithm)
To prevent the Cache Stampede without distributed locks, the XFetch algorithm computes an asynchronous early refresh decision dynamically during read queries:
Where:
- is the time the last recomputation took (stored with the value).
- is the key's expiry time and the current time.
- is the aggressiveness multiplier (typically ).
- is a random floating-point value.
- Because is exponentially distributed, each read has a small chance of refreshing early, and the chance grows as expiry approaches; a key read often is almost certainly refreshed shortly before it expires, and a rarely read key almost never early. It makes a stampede unlikely, not impossible.
Eviction Policies & In-Memory Data Structures
When cache memory reaches its configured maxmemory limit, the eviction policy determines which keys are purged:
Synthesizing vector architecture diagram...
Two structures work together here. The hash map (top) maps each key straight to its node in the list, so a lookup is O(1) with no scanning, as the dotted arrow to Node B shows. The doubly linked list runs from HEAD (most recently used) to TAIL (least recently used). On a hit, the node is unlinked from its place and moved to HEAD, which is O(1) because each node knows both neighbors. When memory is full, the node at TAIL is evicted and its key removed from the map. Neither structure alone is enough: the map cannot track order, and the list cannot find a key without scanning.
In-Memory Node Layout
A textbook (and memcached-style) LRU node inside a doubly-linked list contains pointers and payload metadata. Redis and Valkey don't keep this list: they store a small access clock in each object and evict by sampling a few keys (maxmemory-samples: 5 in open-source Valkey, 3 in ElastiCache's parameter groups) and dropping the least recently used of the sample:
| Node Field | Type / Size | Purpose |
|---|---|---|
prev | uint64_t* (8 Bytes) | Pointer to previous node toward Head |
next | uint64_t* (8 Bytes) | Pointer to next node toward Tail |
key | char* (8 Bytes) | Pointer to cached key string |
val | void* (8 Bytes) | Pointer to serialized payload value |
expires_at | int64_t (8 Bytes) | Unix timestamp for TTL expiration |
| Total Overhead | 40 Bytes | Per-entry node overhead (excluding jemalloc chunk padding) |
State-Transition Trace: LRU Doubly-Linked List Eviction
Under capacity limit entries, with initial state [C, B, A] (Head = C [MRU], Tail = A [LRU]):
| Step # | Event / Input | In-Memory / Distributed State | Evaluation & Transition | Outcome / Cache Latency |
|---|---|---|---|---|
| Initial | Baseline state | Capacity , Full Map: {A, B, C} | Pointers: HEAD <-> [C] <-> [B] <-> [A] <-> TAILHead = C (MRU), Tail = A (LRU) | Cache full ( slots allocated) |
| 1 | GET B (Read Hit) | Hash map lookup finds B pointer in | Unlink B.prev & B.next; re-splice B to HeadOrder becomes [B <-> C <-> A]; Tail A unchanged | Cache HitB promoted to MRU; A remains LRU |
| 2 | SET D = "val_d" (Write Insert) | Key D absent; capacity ceiling reached () | Evict Tail (A): delete map["A"], free node AAllocate node D, splice as new Head | Eviction + Insert Order: [D <-> B <-> C]; C is now LRU Tail |
| 3 | SET C = "v2" (Update Hit) | Key C found at Tail position | Update value pointer C.val = "v2"; unlink from TailSplice C to Head; B becomes new Tail | Update + Promotion Order: [C <-> D <-> B]; B is now LRU Tail |
Eviction Algorithm Comparison Matrix
| Eviction Algorithm | Underlying Data Structure | Time Complexity | Memory per Key | Vulnerability / Failure Mode | Real-World Adoption |
|---|---|---|---|---|---|
| LRU (Least Recently Used) | Doubly-Linked List + Hash Map | get/put | of pointers (list); a few bytes of clock per key in Redis/Valkey | Scan Resistance Failure: A full table batch scan flushes entire working set. | Memcached (per slab class); Redis/Valkey allkeys-lru approximate it by sampling, with no list |
| LFU (Least Frequently Used) | Frequency Buckets + Linked Lists (textbook); an 8-bit logarithmic counter per key in Redis/Valkey | get/put | Textbook: tens of bytes; Redis/Valkey: 8 bits of counter | Historical Bias: in textbook LFU, keys that collected many hits resist eviction long after going cold; Redis/Valkey counters decay over time (lfu-decay-time), which mostly handles this. | Redis/Valkey allkeys-lfu (evicted by sampling) |
| FIFO (First In, First Out) | Queue / Circular Ring Buffer | Sub-optimal hit ratio; evicts hot items simply due to age. | Simple streaming buffers | ||
| Random Eviction | Reservoir / Uniform Sampling | High variance; randomly evicts critical hot keys. | Redis allkeys-random | ||
| W-TinyLFU | SLRU + Count-Min Sketch Filter | Minor CPU overhead on write to update frequency sketch. | Caffeine (Java), Ristretto (Go) |
3. Data Migration, Anti-Entropy & Consistency Protocols
Maintaining consistency between a primary relational database and an external caching tier across distributed networks requires addressing the Dual-Write Concurrency Problem:
Synthesizing vector architecture diagram...
Follow the numbered steps in time order. Both threads update the same user row. (1) Thread 2's UPDATE reaches the database last, so the database ends with name='Bob'. (2) Thread 2 then writes Bob to the cache. (3) Thread 1's cache write, delayed on the network, arrives after that and overwrites the cache with Alice. The result: the database says Bob and the cache says Alice, and nothing will fix it until the key expires. The usual fix is to stop two writers from both setting the cache: delete (or version-stamp) the key after the database write, so the next read reloads the true value. A bare delete can still be undone by a slow refill, which is why fills should compare versions.
The Invalidation Rule: Never Update, Always Invalidate
Writing the new value into the cache on every update lets two writers' cache writes cross. Deleting the key instead avoids that, but it doesn't close every race: a reader that fetched the old row just before the commit, or from a lagging replica just after it, can refill the old value after the delete, and it then stays until the TTL. The robust rule is that an invalidation must leave something a refill can check: a tombstone or the new value carrying the row's version, with every fill stored only if nothing newer is there. At minimum, delete after the commit and keep a TTL as the last bound:
Asynchronous Anti-Entropy via Change Data Capture (CDC)
In microservice architectures, application code may crash between updating the database and evicting the cache. To make sure every committed change produces an invalidation even if the app crashes, production systems decouple cache invalidation from application logic using CDC (Change Data Capture):
Synthesizing vector architecture diagram...
- Transactional Guarantee: The application writes only to the database. The database writes changes atomically to its append-only transaction log (WAL / MySQL Binlog).
- Streaming Invalidation: A CDC connector (e.g., Debezium) streams change events to Apache Kafka.
- Consumer Eviction: Dedicated invalidator workers consume events and execute
DEL {tenant}:user:{id}against Redis clusters, so no invalidation is lost when an app instance crashes. It doesn't stop a fill that read a lagging replica or paused before the delete from storing the old value afterwards; versioned fills or invalidating from the replica's own stream close that.
3. Online Cluster Resharding, Anti-Entropy & Graceful Draining
- Slot Migration Protocol: Redis Cluster distributes 16,384 slots across primary nodes. During dynamic scaling:
- Set target shard state:
CLUSTER SETSLOT <slot> IMPORTING <source_node_id>. - Set source shard state:
CLUSTER SETSLOT <slot> MIGRATING <target_node_id>. - Batch migrate keys atomically:
MIGRATE <target_ip> <port> "" 0 5000 KEYS <key1> <key2>.... - Commit slot assignment: broadcast
CLUSTER SETSLOT <slot> NODE <target_node_id>to all cluster nodes.
- Set target shard state:
- Client Redirection Handling (
ASKvs.MOVED): In-flight keys on migrating slots return-ASK <slot> <target_ip:port>. The client issues anASKINGprefix command to the target node for that single query without updating its local slot routing cache. Once the slot fully commits, queries return-MOVED <slot> <target_ip:port>, prompting clients to refresh their cluster topology map asynchronously. - Replica Synchronization & Draining: Replicas synchronize state via an in-memory ring buffer (
repl-backlog). On brief network blips, partial resynchronization (PSYNC <replication_id> <offset>) avoids full RDB disk dumps. During node decommissioning, operators trigger a plainCLUSTER FAILOVERon the replica: it pauses clients on the primary, waits until its replication offset matches, then takes over, so no acknowledged write is lost.CLUSTER FAILOVER FORCE/TAKEOVERskip that handshake and are reserved for a primary that is already unreachable.
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~46%). Spend 1 Coin to unlock the remaining 7 production deep-dive sections for a full 24 hours.