Caching & Invalidation
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.
Part 0. Start here
The problem: one price change and 12,000 page views a second
A shop's product reads go through a Valkey cache in front of a PostgreSQL read replica. At 20:00:00 a marketing email reaches 2 million people: the kettle k7 is on flash sale, 39. The price is changed in the database at 20:00:00.000. By 20:01, 12,000 page views a second reach k7's page, on top of 20,000 product reads a second for the rest of the catalog. The 12,000 is the number in the drill The Product Page That Melted Redis.
Two numbers make the problem felt.
The load. The replica answers at most 1,000 product queries a second (our example's number). Every read the cache doesn't answer becomes a query:
| Cache hit ratio | Arithmetic at 32,000 reads/s | Replica queries/s | Share of the replica |
|---|---|---|---|
| 99% | 32,000 × 0.01 | 320 | 32% |
| 98% | 32,000 × 0.02 | 640 | 64% |
| 97% | 32,000 × 0.03 | 960 | 96% |
| 0% (the cache is empty) | 32,000 × 1 | 32,000 | 32 times its ceiling |
At 99%, losing one point of hit ratio doubles the database's load, losing three points quadruples it, and an empty cache sends 32 times what the replica can answer.
The freshness. Nothing in the chain is told that the price changed, except by whoever changes it. If the page's origin sends no Cache-Control header, a CloudFront cache policy's default TTL is 86,400 s, one day, and every browser keeps its own copy on top. In our story the app deletes the cache entry 5 ms after the commit, and the page still shows $49 for ten minutes, to about 6.8 million page views. A read that missed 6 ms after the delete refilled the cache from a replica that was 200 ms behind.
The price changed at 20:00:00.000 and the app deleted the cache entry 5 ms later. List every place the old price can still be served from, and for how long. Then: the database was sized for a 99% hit ratio. What happens to it on the day the cache is empty, and what would you build so that day isn't an outage?
The big picture
Synthesizing vector architecture diagram...
What to notice: the read replica follows the primary's log, so it has a subscription and its staleness is its lag. Valkey, the shield, the edges and the browsers have none: the crossed dashed lines are the messages that nobody sends. Each copy is fresh only because it expires (its TTL) or because someone deletes or overwrites it.
What you'll be able to do after this page
- Explain why a cache is load-bearing, turn a hit ratio into database load, and say what a TTL costs in misses (Part 1).
- Compare cache-aside, read-through, write-through, write-behind, refresh-ahead and a stream-fed copy on equal terms (Part 2).
- Explain the race that puts an old value back after a delete, and close it with versions or leases (Part 3).
- Choose where invalidations come from: the app, the change stream, a new key per version, or only the TTL (Part 4).
- Stop stampedes and avalanches with jitter, one fill per key, stale-while-revalidate and early refresh (Part 5).
- Cache "absent" correctly, and say why it doesn't stop enumeration attacks (Part 6).
- Size a cache from its working set and pick an eviction policy (Part 7).
- Survive a dead cache node and a cold restart without taking the database down (Part 8).
- Give a writer read-your-writes through a cache, invalidate a second Region safely, and use DAX without surprises (Part 9).
- Set
Cache-Control, build a cache key, add up the staleness of every HTTP tier, and choose between a purge and a versioned URL (Part 10). - Trace one page view and one price change through every layer (Part 11).
- Map all of it to AWS services, and name the look-alikes that aren't caches (Part 12).
You may have arrived from a step that relies on this: step 1.5 of the URL shortener loop (a disable overwrites the entry, and refills use SET NX), step 2.2 of the news feed loop (the DAX query cache that writes don't invalidate), step 2.1 of the hotel reservation loop (a change-stream-fed cache that applies only newer versions) or step 1.1 of the Shopify case study (a version in the cache key, so nothing is deleted key by key). This page is the "why" behind all four, and behind about thirty other loop steps that put a cache in the path.
Part 1. What a cache is, and why the database now depends on it
The catalog gets 20,000 product reads a second, and the replica takes 1,000. Without a copy of the answers somewhere cheaper, the shop can't open. So we keep one, and from that moment the database's health depends on how often that copy answers.
A copy with no subscription
A cache holds copies of answers the database already gave. The database doesn't know the copies exist, so it can't tell them when a row changes. A copy becomes fresh again in only two ways:
- it expires: its TTL (time to live) runs out, and the next read fetches a new copy;
- someone invalidates it: deletes it or overwrites it after a write.
A few more words used on this page:
| Word | Meaning |
|---|---|
| Hit / miss | The cache has a usable copy / it doesn't |
| Fill (or refill) | On a miss, read the source and store the answer in the cache |
| Eviction | The cache drops an entry to make room, whatever its TTL |
| Working set | The keys that are read often enough to be worth keeping |
| Hit ratio | Hits ÷ reads over some window |
The read path and the write path
The default design is cache-aside: the application talks to the cache and to the database itself.
textREAD(k) (cache-aside) 1. v = GET product:{k} 2. hit -> return v 3. miss -> v = query the source (here: the read replica) SET product:{k} v EX 600 return v WRITE(k, new value) 1. commit the change in the database (the primary) 2. THEN invalidate product:{k} (never before the commit; Part 3 decides what "invalidate" must write)
A miss costs three trips: the cache, the database, and the cache again.
Hit ratio and the database
The database only sees the misses. Formula 1:
That is the hook's table. Because the database is sized for the misses, a small drop in hit ratio is a large rise in its load: at 32,000 reads a second, 99% means 320 queries, 97% means 960, and 96% means 1,280, already over the replica's 1,000.
The replica takes 1,000 queries a second and we serve 32,000 reads. What hit ratio must the cache hold, and what happens at 3 points below it?
What a TTL costs in misses
A key read at rate λ is cached for one TTL after each miss, and then waits on average 1/λ for its next read, which misses. So it misses once every TTL + 1/λ seconds. In words: each key that is read costs at most one miss per TTL, and a hot key (many reads per TTL) costs almost exactly one. So the misses a second are at most the number of distinct keys read in one TTL ÷ the TTL.
For our catalog (event 1, from the reference implementation):
- The bound: 100,000 keys ÷ 600 s = 167 misses a second at most.
- The exact steady state, summing each key's 1 ÷ (TTL + 1/λ): 159 misses a second, a hit ratio of 99.2%, 16% of the replica. Cold products read less than once per TTL cost less than one miss per TTL each.
- At 19:59 the 10,000 most popular products were loaded by a job at 19:55 and none of them is due before 20:04, so only the others miss: 142 a second (14% of the replica).
Halving the TTL roughly doubles the misses of the keys that are read often; it can't more than double them.
Why not more read replicas instead?
Each replica adds about 1,000 queries a second, and each must apply every write (the Replication, Quorums & Read-Your-Writes loop primitive). At 32,000 reads a second with no cache, that is 32 replicas. With a 99% cache it's one. A cache turns a read-capacity problem into a memory problem: keep the working set in memory. The price is the rest of this page: staleness you must bound, and a database that can't survive the day the cache is empty.
Why not cache it forever?
The drill's second question asks exactly this. "Cache the whole catalog forever" means a changed price is invisible until something deletes the entry, and Part 3 shows how a delete can be undone by a racing refill. The TTL is the last bound on every invalidation you missed. With no TTL, a missed one is never repaired.
The example we follow
Our shop is small enough that each layer's arithmetic fits on a line; the numbers are smaller than a big shop's. Every trace on this page comes from running a private reference implementation of this shop, not from working by hand.
| Setting | Tiny example | At real scale |
|---|---|---|
| Catalog | 100,000 products k1 … k100000; kN is the N-th most popular. Popularity follows a Zipf curve with exponent 1.0 (our choice): the top 10,000 get 81.0% of reads, the top 20,000 86.7%, the top 50,000 94.3%. k7 alone gets 236 reads a second before the sale | Millions of products; measure the skew, don't assume it |
| Product row | products(id, price, stock, description, version); version goes up by 1 on every write to the row | A per-row counter or the commit's log position (Change Streams & the Transactional Outbox) |
| The price change | 20:00:00.000: k7 39 v8**, one UPDATE on the primary | |
| Traffic at the app tier | 20,000 catalog product reads a second. From 20:00:05, k7 ramps to 12,000 page views a second by 20:01:00. Each view makes one public request /p/k7 (the page with the price, cacheable at the CDN) and one private request /api/me/k7 (member price, "in your cart"; never cached publicly) that reads product:{k7} | The same |
| Layers | Browser → CDN (two edge locations, E1 and E2, and one Origin Shield) → load balancer → 8 app servers (A1 to A8) → Valkey (cache-aside) → PostgreSQL read replica for fills, primary for writes and checkout | CloudFront has "750+ POPs" and 15 regional edge caches (AWS's figures); tens of app servers |
| App server capacity (labelled) | 5,000 requests a second each: 40,000 for the tier | Load-test |
| Valkey cluster | 4 shards, A to D, over 16,384 slots in equal ranges (A 0 to 4,095, B 4,096 to 8,191, C 8,192 to 12,287, D 12,288 to 16,383). Keys hash on the tag in {}: product:{k7} → slot 4,452, shard B. Shard B holds 22.7% of catalog reads, 4,539 a second. Each shard is a primary only: the team removed the replicas last month to save money ("it's only a cache"). Memory holds the whole catalog, about 100,000 × 2 KB ≈ 200 MB | ElastiCache for Valkey, cluster mode, 0 to 5 replicas per shard |
| Cache entry | product:{k7} = {price, stock, description, version}, TTL 600 s, no jitter; the 19:55 pre-warm job already uses ±10% jitter | Loops use 1 s to 1 h |
| Fill and invalidation (naive until Part 3) | Miss → read the replica → SET … EX 600. A write: commit, then DEL product:{k} | |
| Database | One read replica for fills: 1,000 product queries a second (20 connections × 20 ms per query, labelled). Replica lag during the sale: 200 ms (labelled). The primary also answers in 20 ms | Measure; lag is page 06's subject |
| Overload model | A FIFO queue in front of the 20 connections. A queued query whose reader has given up is dropped before it runs (the pool's acquire timeout, or the request's cancellation). A fill whose reader timed out does not write the cache. The app times out after 1 s and retries once | The connection pool's acquire timeout; request cancellation; driver timeouts |
| CDC invalidator (from Part 4) | Tails the primary's change stream, writes a marker per changed row; lag about 50 ms (labelled) | Debezium, AWS DMS, DynamoDB Streams to Lambda |
| Negative cache (from Part 6) | "absent" entries with a TTL of 60 s | Loops: 30 s to 5 min |
| Write-behind counter | views:{k7} counted in Valkey on every view, flushed to the database every 10 s | The ad-click and YouTube loops |
CDN settings for /p/k7 | The origin sends Cache-Control: public, max-age=30, s-maxage=5, stale-while-revalidate=10, stale-if-error=300 (our values). The cache policy keys on the path only (no query strings, no cookies), minimum TTL 0 | CloudFront cache policies |
| Second Region (Part 9) | eu: its own app servers and Valkey, a cross-Region replica 1 s behind (labelled); invalidations broadcast from home reach eu in 80 ms (labelled) | Aurora Global Database, DynamoDB global tables |
| The merchandiser (Part 9) | Sets the price at 20:00:00.000 and reloads the product page at 20:00:02 |
The overload model decides the overload numbers. Whether a query whose reader has gone still runs, and whether a late answer still writes the cache, changes the result completely. Every overload number on this page (Parts 5 and 8) uses the model in the table.
Snapshot S1, at 19:59:10, from the reference implementation:
Copy of k7 | Value | Until |
|---|---|---|
| Primary | $49 v7 | |
| Read replica | $49 v7 (it follows the primary, 200 ms behind) | |
| Valkey shard B | $49 v7 | 20:04:07.0 (filled 19:55:00 by the pre-warm job, TTL 547 s after jitter) |
| Origin Shield, edges E1 and E2 | $49 | 5 s after each one's last fetch |
| A viewer's browser | $49 | 30 s after its last load |
eu Valkey and eu replica | $49 v7 | |
| Replica load | 142 queries/s, 14% of 1,000 |
The story has twelve beats, each in its own Part: the cache is load-bearing (Part 1); the same price change under five patterns, on a copy (2); the price changes and the old price comes back (3); an invalidation gets lost (4); a batch of keys and then k7 expire under load (5); a deleted product keeps being asked for (6); the cache shrinks, on a copy (7); a cache node dies (8); the merchandiser and the second Region (9); back to 20:00 at the edge (10); one page view and one price change, end to end (11); and on AWS (12). Replays marked "on a copy" run a what-if without changing the main timeline. The full event table is in Part 13.
Background reading, not relied on for any fact here: Primitive #04: Distributed cache patterns and eviction.
What to remember from Part 1
- A cache is a copy nobody tells about changes: it's fresh only by expiry or by invalidation.
- Database load = reads × (1 − hit ratio): at 99%, one point of hit ratio doubles it.
- Each key that is read costs at most one miss per TTL; a hot key costs almost exactly one.
Part 2. Five ways to put a cache in the path, and one alternative
The price changes. Who updates the cache: the code that writes the database, a library, the cache itself, or nobody? Where the cache sits in the write path decides what can go wrong, so we replay the 20:00:00 change on a copy under each design.
Two writers cross
Take write-through: every write goes to the database, then sets the cache to the new value. It sounds like "the cache is never stale". Now two writers change k7 at nearly the same moment: the merchandiser's 38 (v9) ten milliseconds later. Writer 1 is slowed down on its way to the cache.
| Time | Writer 1 (merchandiser) | Writer 2 (correction) | Database | Cache |
|---|---|---|---|---|
| .000 | Commits $39 v8 | $39 v8 | $49 v7 | |
| .010 | (paused) | Commits $38 v9 | $38 v9 | $49 v7 |
| .012 | Sets the cache: $38 v9 | $38 v9 | $38 v9 | |
| .018 | Sets the cache: $39 v8 | $38 v9 | $39 v8 |
Side row P2, from the reference implementation. The database says 39 until the TTL, at 20:10:00.018. With a version compare on the cache write ("only if the stored version is lower"), writer 1's v8 is refused and the cache keeps $38 v9.
Synthesizing vector architecture diagram...
What to notice: in "Cache-aside" and "Write-through" the database is written first and stays the truth; the cache follows. In "Write-behind" the order is reversed: until step 2 the cache is the only place that knows the new price.
Write-through updates the cache on every write. Why can the cache still end up older than the database?
Six designs on equal terms
| Design | When a reader sees the new price | A crash between the two writes leaves | What a miss means | Write latency | Cached but never read | Where the truth lives |
|---|---|---|---|---|---|---|
| Cache-aside (app fills on a miss, deletes after a write) | After the delete, on the next fill (if the fill is ordered: Part 3) | The old value, until the TTL | "Not cached": read the source | Commit + one delete | Nothing: only keys someone read | Database |
| Read-through (a library or the cache fills on a miss) | Same as cache-aside: the same races | Same | Same | Same | Nothing | Database |
| Write-through (write the database, then set the cache) | Right after the write, if writes don't cross (P2) | The old value, until the TTL | "Not cached": read the source (cold keys still miss once) | Commit + one set | Every written key, read or not | Database |
| Write-behind (write the cache, flush later) | Right away, from the cache; the database only after the flush | A lost write if the cache node dies first | "Not cached": but the database may be behind | One cache write | Every written key | The cache, until the flush |
| Refresh-ahead (reload before expiry) | Does nothing for a write: only for keys about to expire | Not applicable | As cache-aside | Unchanged | Keys reloaded that nobody asks for again | Database |
| A complete copy fed by the change stream (a view, not a cache) | After the stream's lag, in order, applied by version | Nothing: the stream is retried until applied | "Doesn't exist": every row is present, so absence is an answer | Unchanged (the stream does the work) | Every row | Database; the copy is derived and rebuildable |
Check the rows from both sides. Write-through doesn't help a cold key's first read (nobody wrote it), and costs memory for keys nobody reads. Cache-aside adds almost nothing to write latency (one delete), and keeps the cache small. Refresh-ahead helps reads of hot keys (Part 5 makes it self-tuning with XFetch), not freshness after a write.
Write-behind: the cache is the truth until it flushes
Write-behind turns the cache into the source of truth until the next flush. Replayed for the price (side row P3): until the flush, the database still says 49; and if the cache node is lost before the flush, the price change is lost with it. So we use write-behind only for data we can lose or rebuild. In our shop that's the "trending" counter views:{k7}: one increment per view in Valkey, flushed every 10 s (Part 8 counts what a node failure loses).
A complete copy fed by the stream
Several loops don't use a cache at all where a missing key would be ambiguous. They keep a complete copy that an updater writes from the change stream, every row present, applied only if newer: the proximity service (step 2.1), the digital wallet (step 2.4) and ad-click aggregation (step 2.1). Its lag is the stream's lag, and a missing key means "doesn't exist", which a cache can never say. The updater itself is the cache updater of the Change Streams & the Transactional Outbox loop primitive. DAX, AWS's cache for DynamoDB, is a managed write-through item cache; Part 9 has its details.
What to remember from Part 2
- Cache-aside is the default: the database stays the truth, and a cache failure only costs misses.
- Write-through still needs versions when two writers race.
- Write-behind makes the cache the truth until it flushes: only for data you can lose or rebuild.
Part 3. The race that leaves a stale value for good
The price changed at 20:00:00.000, and the app deleted the cache entry 5 ms later. Yet the page shows $49 for ten minutes. A delete leaves nothing behind, so a refill that read the old price can land after it, and nothing stops it.
Trace: the old price comes back
From the reference implementation (the naive design: DEL after commit, plain SET on refill; k7 still at its pre-sale 236 reads a second):
| # | Time | Event | Valkey product:{k7} |
|---|---|---|---|
| 3 | 20:00:00.000 | The primary commits k7 $39 v8. The replica will apply it at .200 | $49 v7 (from the pre-warm) |
| 4 | 20:00:00.005 | The app runs DEL product:{k7} | empty |
| 5 | 20:00:00.011 | A1 reads k7: a miss. It queries the replica, which still has v7 | empty |
| 5 | .019 to .024 | A4, A3 and A8 miss too, and query the replica | empty |
| 5 | 20:00:00.032 | A1's fill lands: SET … $49 v7 EX 600 | $49 v7 |
| 5 | .040, .044, .045 | The other three fills land, each a plain SET that restarts the TTL. The last one, at 20:00:00.045, sets the expiry: call it t₅ | $49 v7 until 20:10:00.045 |
| 6 | 20:00:05 to 20:10:00.045 | Every read of k7 hits 49 to about 6.8 million page views | $49 v7 |
Checkout still charges $39, because it reads the primary. That's the rule the drill's gold answer depends on: the cache decides what is shown, never what is charged. Money, stock and permissions are decided against the source.
Snapshot S2, after event 4:
| Copy | Value |
|---|---|
| Primary | $39 v8 |
| Read replica | $49 v7 (applies v8 at .200) |
| Valkey shard B | empty (deleted) |
Synthesizing vector architecture diagram...
Snapshot S3, from the reference implementation. What to notice: the delete comes first and the stale SET second. At .200 the replica is right and the cache is still wrong, and nothing will tell it. (Three more fills repeat the same SET by .045.)
The app deleted the entry 5 ms after the commit. How did $49 get back in?
Race B: the lagging source
What we just saw. A fill that starts after the delete reads a copy that is behind (a replica, or another Region's copy), so it gets the old value and stores it after the delete. Reading the primary instead would close this race.
Race A: the slow reader
A fill that read the primary can be stale too, if it read before the commit and stored after the delete. Side row 6a, on a copy where k7's entry had expired at 19:59:59.985:
| Time | Event | Valkey |
|---|---|---|
| 19:59:59.990 | A5 misses and queries the primary; the query reads the row as it is when it starts: $49 v7 | empty |
| 20:00:00.000 | The primary commits $39 v8 | empty |
| 20:00:00.005 | The app deletes the key (nothing to delete) | empty |
| 20:00:00.010 | A5's query returns $49 v7; then a 40 ms garbage-collection pause | empty |
| 20:00:00.050 | A5 stores $49 v7 EX 600 | $49 v7 until 20:10:00.050 |
Reading the primary closes race B, not race A.
What the refill can see
Three ways to write the invalidation and the refill, each replayed on both races (side row 6b, from the reference implementation):
| Invalidation / refill | Race A (A5 at .050) | Race B (A1 at .032) | Why |
|---|---|---|---|
(i) DEL, then refill with SET … NX (only if absent) | Stored: $49 v7 | Stored: $49 v7 | After a delete the key is absent, so NX succeeds with the old value. NX is "first writer wins", whatever its age |
(ii) Overwrite with the new value, refill with a plain SET | Stored: **39 | No miss: the key holds $39 | A slow refill overwrites the new value |
(iii) Overwrite with the new value, refill with SET … NX | Refused: $39 v8 stays | No miss: the key holds $39 | The refill finds the new value and does nothing. This is the URL shortener loop's rule (step 1.5) |
(iii) closes both races for one writer. Two writers can still cross, exactly as in P2: their overwrites arrive out of order, unless the overwrite compares versions. SET NX protects a value the write path wrote, not a key the write path deleted.
Two more ways people try, both side rows in Part 13:
- Delete before the commit (6c): the delete at 19:59:59.995 lets a reader miss at .997 and read v7 before the commit; its SET lands at 20:00:00.018, after the commit: stale until the TTL. Always invalidate after the commit.
- The delayed double delete (6d): delete again 500 ms later. It fixes event 5: the stale fills landed by .045, the second delete at .505 removes them, and the next fill reads the replica, which has had v8 since .200. It doesn't fix a reader paused for 1 s: its SET lands at 1.010, after the second delete. It narrows the window by a guess; it doesn't close it. The news feed loop uses it (step 3.2) with a 60 s TTL behind it.
Why does anyone delete at all? The memcache paper from Facebook says it plainly: "We choose to delete cached data instead of updating it because deletes are idempotent." A delete can be repeated and reordered safely. The paper then adds leases to order the refills, which is Fix 2 below.
Fix 1: versions on both sides
Every row already has a version from one authority, the database. The invalidation writes a tombstone carrying the new version. Both the invalidation and every fill are compare-and-sets, so nothing older ever replaces something newer:
textINVALIDATE(k, v) after the commit returned version v; one server-side script, atomic cur = product:{k} if cur holds version >= v -> do nothing (a newer write already got here) else write TOMB(v) with a TTL of 60 s (60 s > the slowest fill + the worst replica lag a fill can read) or, equally: write (new value, v) under the same compare FILL(k, value, v') one server-side script, atomic cur = product:{k} if cur is empty -> store (value, v'), TTL 600 if cur is TOMB(t) and t <= v' -> store (value, v') if cur is a value with version < v' -> store (value, v') else -> refuse (a refused filler re-reads the primary, or serves its value once, marked stale)
Valkey's SET has NX, XX, IFEQ and IFNE (set only if the current value equals, or doesn't equal, a given one), but no "only if newer". A version compare needs a short server-side script. The nearby friends loop (step 1.2) keeps the version in a separate key and uses WATCH/MULTI/EXEC: a refill writes only if the version didn't change while it ran. That stops the slow reader, but not a refill that starts after the change and reads a lagging copy, so those refills must read the source.
Event 7 replays events 3 to 6 with Fix 1 (from the reference implementation):
| Time | Event | Valkey |
|---|---|---|
| .000 | Commit $39 v8 | $49 v7 |
| .005 | The app writes TOMB v8 (TTL 60 s) | TOMB v8 |
| .011 | A1 misses (a tombstone reads as a miss) and queries the replica: $49 v7 | TOMB v8 |
| .032 | A1's fill: refused, 7 < 8. A1 re-reads the primary | TOMB v8 |
| .053 | A1 stores $39 v8, 21 ms after the refusal | $39 v8 |
Twelve fills were refused in all (every reader that missed before .053), each followed by a primary read: 24 queries instead of 4. The last read that got 39 v8 at .000 and writer 2 commits $38 v9 at .010; writer 2's TOMB v9 lands at .012 and writer 1's TOMB v8, delayed, at .018; at .205 a reader reads the replica, which has v8 but not yet v9 (it applies v9 at .210), and its fill lands at .226.
| Tombstone write | After .018 the key holds | The v8 fill at .226 | Result |
|---|---|---|---|
| Unconditional | TOMB v8: the older marker replaced the newer one | Stored (8 ≤ 8) | **38 v9 |
| Compare-and-set | TOMB v9: the v8 write was a no-op | Refused (8 < 9) | The filler re-reads the primary and stores $38 v9 |
TOMB v8 is this page's invalidation marker: "version 8 exists; don't store anything older". The Change Streams & the Transactional Outbox loop primitive (Part 8) stores GONE@6 for a deleted row. Same shape, different meaning, and the same lesson: caches must store something, so a late refill can't bring the old value back.
Fix 2: leases
The memcache paper introduced leases "to address two problems: stale sets and thundering herds". On a miss, the cache hands the client a 64-bit token bound to the key. A delete of the key voids its tokens, and a SET with a void token is refused.
| # | Replay | Result, from the reference implementation |
|---|---|---|
| 7a | Race A with a lease | A5's token was issued at 19:59:59.990 and voided by the delete at .005: its SET is refused. Fixed |
| 7a | Race B with a lease | A1's token is issued at .011, after the delete, so it's valid: $49 v7 is stored |
Leases close the slow reader, not the lagging source. The paper says as much for its own replicas: "A cache refill from a replica's database should only be allowed after the replication stream has caught up." So with leases, fill from the primary, or invalidate only after the copy you fill from has the write (Part 4). (A lease here is a cache mechanism, not the ownership lease of the Leases, Fencing Tokens & Distributed Locks loop primitive.)
The TTL is the last bound
With no TTL at all (side row 7b), event 5's $49 v7 would stay until an eviction or the next price change. Every fix above has rare holes; the TTL is what bounds them.
Every write-path option on equal terms
| Option | Race A (slow reader) | Race B (lagging source) | Two writers | A lost invalidation | Marker evicted before a stale refill lands | Cost |
|---|---|---|---|---|---|---|
DEL + plain refill | Stale until TTL | Stale until TTL | Fine (deletes don't carry values) | Stale until TTL | Nothing to evict: already open | None |
DEL + SET NX refill | Stale | Stale | Fine | Stale until TTL | Already open | None |
Overwrite + SET NX refill | Closed | Closed | Stale until TTL | Stale until TTL | Open (the new value evicted, then NX stores the old) | One writer only |
| Versioned compare-and-set (tombstone or new value; the invalidation compares too) | Closed | Closed | Closed (7c) | Stale until TTL | Open (the tombstone evicted) | A version column; a script per fill; tombstone memory; extra primary reads |
Version key + WATCH | Closed | Stale (a fill that starts after the bump and reads a lagging copy sees an unchanged version): read the primary, or use a strongly consistent read | Closed | Stale until TTL | Open (the version key evicted) | A second key; a transaction per fill |
| Lease | Closed | Stale | Fine | Stale until TTL | Not applicable | A cache that supports it; readers wait |
| Delayed double delete | Closed only for pauses shorter than the delay | Closed only if lag + fill < the delay | Fine | Stale until TTL | Not applicable | A timer per write; a guess |
| TTL only | Stale up to the TTL | Stale up to the TTL | Stale up to the TTL | Nothing to lose | Not applicable | Staleness = the TTL |
The eviction column needs no event: a marker is evicted between the write and a stale refill only under memory pressure, and the TTL bounds it. A lost invalidation is Part 4's subject.
What to remember from Part 3
- Deleting on write isn't enough: a refill that read the old value can land after the delete.
- An invalidation must leave something the refill can check: the new value (refill with
NX) or a version (refill only if newer). - A lease stops a slow reader, not a lagging replica; the TTL bounds whatever you missed.
Part 4. Where invalidations come from
Fix 1 needs someone to write the tombstone after every commit. From 20:02:00 the team turns Fix 1 on by a flag (the stale k7 entry from event 5 is already stored and stays). One minute later, the app's own invalidation gets lost.
Trace: the invalidation that timed out
k30003 is on shard B too (slot 5,231). It's the 30,003rd most popular product: one read every 18 seconds, and it wasn't pre-warmed. From the reference implementation:
| # | Time | Event | Valkey product:{k30003} |
|---|---|---|---|
| 20:02:46.382 | A read misses and fills v4 | v4 until 20:12:46.382 | |
| 8 | 20:03:00.000 | The primary commits k30003 v4 → v5 | v4 |
| 8 | 20:03:00.005 | The app's TOMB v5 write to shard B times out (a network blip). Nothing retries it | v4 |
| 20:03:17 to 20:12:31 | 35 reads get the old price | v4 | |
| 20:12:46.382 | The entry expires | empty | |
| 20:13:06.797 | The next read (at .776) misses and fills v5 from the replica | v5 |
The old price was served for 607 s: the rest of the entry's TTL, plus the wait for the next read. The app's invalidation is a second write after the commit. A timeout, a crash or a deploy between the two loses it, and nothing notices.
Sources on equal terms
| Source | Staleness | What a crash or timeout loses | Cost | Cold keys it creates |
|---|---|---|---|---|
| The app, after its commit | Milliseconds | The invalidation, silently: stale until the TTL | Nothing extra | None |
| The change stream (a CDC invalidator reads the database's log) | The invalidator's lag (50 ms here) | Nothing: it retries until the cache acknowledges, in order per key | Stream plumbing and an invalidator to run and alarm on | None |
A new key per version (page:{k7}:v8) | The pointer's TTL (below) | Nothing to lose: nothing is invalidated | A pointer read per request; memory for old versions until they age out | Every version bump: a bump of a whole shop is a cold cache |
| TTL only | Up to the TTL | Nothing to lose | Nothing | One miss per key per TTL (Part 1) |
| In-process copies on 8 app servers | A broadcast's delay, or a TTL short enough to accept (5 s here) | A broadcast missed by a restarting server: until the TTL | A pub/sub channel, or server-assisted client tracking (Part 11) | None |
| Combined: app fast path + stream guarantee + TTL bound | Milliseconds usually; the invalidator's lag at worst; the TTL if both fail | Only a double failure, bounded by the TTL | The sum | None |
We run the combination. The app's tombstone stays as the fast path; the stream is the guarantee; the TTL is the bound.
Synthesizing vector architecture diagram...
What to notice: three paths write the same marker. The dashed app path is fast but can be lost; the primary's stream is retried until it lands; the replica's stream acts only after the copy that fills read has the write. The TTL bounds whatever all three miss.
| # | Replay | Result, from the reference implementation |
|---|---|---|
| 8a | Event 8 with the CDC invalidator | At 20:03:00.050 it reads v5 from the change stream and writes TOMB v5 (acknowledged at .051). The next read, at 20:03:17.822, misses and fills v5 from the replica, which has had it since .200. Had a read come before .200, its fill of v4 would have been refused (4 < 5) |
A repeated tombstone is harmless: it's an absolute write, not an increment, so applying it twice leaves the same state, and the compare turns an older, replayed tombstone into a no-op (the Idempotency & Effectively-Once Processing loop primitive, Part 5). From 20:04:00 the CDC invalidator runs in the main line. It and the app's fast path write through the same compare-and-set script, so whichever lands second, or a replay hours later, can't put an older version over a newer one.
Your invalidator reads the primary's change stream and your fills read a replica. What can still go wrong, and which of two changes fixes it?
Invalidate from the copy you fill from
An invalidator that tails the replica's change stream acts only after the replica has applied the change. Replayed for k7 at 20:00 with no versions at all (side row 8b, from the reference implementation): the replica applies v8 at .200 and its invalidator deletes the key at .250. Event 5's stale fills (landed by .045) are deleted then, and the next miss reads the replica, which has v8: $39 v8 is stored. Race B is closed by ordering, without versions.
Race A is not. A reader that read the replica at .190 (still v7), paused 70 ms and stored at .260 lands after the .250 delete. The versioned fill still closes that one. This ordering rule is the memcache paper's own ("A cache refill from a replica's database should only be allowed after the replication stream has caught up"), and it's why the Figma case study's invalidator tails one replication stream per shard.
A new key per version
Some caches avoid invalidation entirely: the key names the version. The rendered page for k7 is cached as page:{k7}:v7; readers learn the current version from a small pointer ver:{k7} that they cache for 1 s (labelled). Side row 8d, from the reference implementation:
| Time | What readers do |
|---|---|
| .000 | Commit v8. Nothing is deleted |
| .000 to 1.000 | Readers whose pointer copy still says 7 keep asking for page:{k7}:v7: the old page, for at most the pointer's 1 s |
| After that | They ask for page:{k7}:v8: a miss, one fill (with a lease), then hits |
| Later | page:{k7}:v7 ages out by TTL or eviction |
There is nothing to invalidate and no refill race, since a refill of v7 writes a key nobody asks for any more. The cost: every bump is a cold key, a pointer read on every request, and the pointer is itself a cached value with a staleness bound. The Shopify case study (step 1.1) puts a shop version in every page key, so a shop-wide edit is a cold cache for that shop (Part 8); notifications (step 1.5) key templates by version and never invalidate them.
The invalidator's lag is your staleness
Side row 8c, on a copy: a bulk import changes 50,000 prices at 20:04:00 in 20 s, 2,500 a second. The invalidator writes 1,000 tombstones a second (labelled), so a backlog builds, and the last change waits 50,000 ÷ 1,000 − 20 = 30 s for its tombstone. While the lag lasts, it is the staleness; the TTL still bounds it. So alarm on the invalidator's lag, and keep the tombstone TTL (60 s) longer than the slowest fill plus the worst replica lag a fill can read. The stream itself (ordering, retention, poison records) is the Change Streams & the Transactional Outbox loop primitive's subject, Parts 3, 4 and 6.
What to remember from Part 4
- An invalidation sent by the app can be lost; one read from the change stream is retried until it lands.
- Invalidate after the copy you fill from has the write, or compare versions.
- The invalidator's lag is your staleness while it lasts.
Part 5. Stampedes: when a hot key or a batch of keys expires
A cache protects the database only while it answers. Misses don't come evenly: they come in herds. Many keys can expire in the same second, or one very hot key can expire under thousands of readers. First, a herd that the team saw on a copy, which is why the main line's pre-warm job uses jitter.
Many keys at once: the avalanche, on a copy
Side replay 14: the same 19:55 pre-warm job without jitter. All 10,000 keys get TTL 600 and expire together at 20:05:00. They carry 81% of catalog reads, so at 20:05:00 misses arrive at about 16,200 a second (0.81 × 20,000) against a replica that answers 1,000. The hottest keys are refilled in the first few tens of milliseconds, so the first second counts 8,243 misses. Under the overload model (a FIFO queue, queued queries whose reader has given up dropped before they run, late answers don't fill, one retry), from the reference implementation:
| Second | Misses | Queries sent (with retries) | Answered in time | Answered late | Cancelled unrun | Timeouts | Failed after the retry | Queue |
|---|---|---|---|---|---|---|---|---|
| 20:05:00 | 8,243 | 8,390 | 985 | 0 | 0 | 0 | 0 | 6,818 |
| 20:05:01 | 11,648 | 11,928 | 80 | 920 | 6,362 | 7,330 | 0 | 11,400 |
| 20:05:02 | 12,134 | 12,426 | 0 | 1,000 | 10,914 | 11,928 | 5,731 | 12,554 |
| 20:05:30 | 12,383 | 12,648 | 0 | 1,000 | 11,636 | 12,617 | 6,371 | 12,760 |
| 20:05:59 | 12,055 | 12,311 | 0 | 1,000 | 11,433 | 12,440 | 6,230 | 12,524 |
The replica keeps working at 1,000 queries a second, but from the third second every answer arrives after its reader gave up. The FIFO queue always holds about a second of requests, so the query at its head is always about to time out: it runs, finishes late, and its fill is thrown away. Over the minute after 20:05:01 the replica ran 59,000 queries and only 80 were answered in time. Just 644 of the 9,999 keys got refilled, so misses never fall: useful work collapses to near zero and stays there until something sheds load. That mechanism (timeouts, retries, shedding) is the Retries, Timeouts, Backpressure & Load Shedding loop primitive's subject (coming). This page keeps the cache's side: cap the fills, fail fast, serve stale.
Two words for two herds:
- Avalanche: many keys, one moment (a batch filled together expires together, or a node dies).
- Stampede (thundering herd): one key, many readers at the moment it expires.
Jitter, the main line
The main line's pre-warm (event 13) gave each key a TTL of 600 s ± 10%, 540 to 660 s. Their expiries spread from 20:04:00 to 20:06:00: 10,000 ÷ 120 s ≈ 83 extra fills a second. From the reference implementation, the replica averaged 225 queries a second (at most 263) over those two minutes, with no queue and no timeouts. Jitter spreads many keys expiring together. It does nothing for one hot key: that key still expires, once, under all its readers.
Trace: k7expires under 12,000 reads a second
Event 5's stale entry expires at t₅ + 600 s = 20:10:00.045. k7 now gets 12,000 reads a second, and nothing collapses them (the protection flag comes at 20:10:05). From the reference implementation:
| # | Time | Event |
|---|---|---|
| 15 | 20:10:00.045 | k7's entry expires. Every read misses, and each sends its own replica query |
| 15 | +21 ms | The first fill lands (20 ms query + 1 ms SET): $39 v8, stored because the key is absent. Reads now hit |
| 15 | In those 21 ms: 12,000 × 0.021 ≈ 250 expected; 231 in this run. With about 3 connections busy with other fills, the queue peaks at 202 | |
| 15 | +270 ms | The queue is empty again. The readers who missed waited up to 220 ms (p99); no timeouts |
The stampede is also the moment the price finally becomes right: 49.
Synthesizing vector architecture diagram...
Synthesizing vector architecture diagram...
Snapshot S4, from the reference implementation. What to notice: about 230 queries are sent in the first 21 ms and the queue then drains at 1,000 a second for about 270 ms. With one fill per key the same expiry costs 8 queries (one per server) or 1.
k7 expires under 12,000 reads a second and each fill takes 20 ms. How many queries does the replica see with no protection, with single-flight on 8 servers, and with a lease?
One fill per key
| # | Replay of event 15 | Result, from the reference implementation |
|---|---|---|
| 15a | Single-flight per app server | 8 queries, one per server; 223 readers waited, up to 20 ms |
| 15b | Lease: the cache grants one token per key per 10 s; others wait 5 ms and try the cache again | 1 query; 230 readers waited, up to 25 ms |
Single-flight and the micro-cache are the Sharding, Hot Keys & Rebalancing loop primitive's tools for one hot key (Part 4): one fetch per process, everyone else waits for it, and a short cache in front. Here they are one row of the comparison.
A lease (Part 3's mechanism) is single-flight across the whole cluster. The memcache paper's servers "return a token only once every 10 seconds per key"; a client refused a token waits a short time and retries, and "typically, the client with the lease will have successfully set the data within a few milliseconds". The paper reports the effect on keys prone to herds: "Without leases, all of the cache misses resulted in a peak database query rate of 17K/s. With leases, the peak database query rate was 1.3K/s." One risk: a lease holder that dies makes the others wait up to the 10 s token interval, unless they can take a stale value.
Serve stale, refresh once
Stale-while-revalidate keeps an entry past its soft expiry: readers get the stored value while one refresh runs.
textREAD(k) main line from 20:10:05: leases + stale-while-revalidate 1. e = GET product:{k} (value, version, soft expiry; hard TTL = soft + 60 s) 2. e is a value and now < e.soft_expiry -> serve e 3. e is a value (hard TTL not up) -> serve e; if we get the lease, refresh in the background 4. no e (a tombstone counts as no e): we get the lease -> fill (compare-and-set, Fix 1), serve someone else holds it -> wait 5 ms, go to 1 (until the request's own timeout)
| # | Replay of event 15 | Result, from the reference implementation |
|---|---|---|
| 15c | Soft TTL 600 s, hard TTL 660 s | 1 query, nobody waits. The stored value is still event 5's 49 during the 21 ms refresh: the price of never making a reader wait |
The memcache paper does the same after a delete: it hands out the deleted value "marked as stale" to readers that can use it.
Refresh before it expires
A hot key could be refreshed a little before it expires, by one reader, while everyone else keeps hitting. XFetch (Vattani, Chierichetti and Lowenstein, "Optimal Probabilistic Cache Stampede Prevention") does it with no scheduler and no coordination. Each read recomputes early if Formula 2 holds:
where Δ is how long the last recompute took (20 ms here), β = 1 is the paper's default, and rand() is uniform in (0, 1], so −ln(rand()) is 0 or more, usually around 1, rarely much larger. Each read pretends it's a little later than it is. A key read 12,000 times a second has many reads in its last 100 ms, and one of them fires; a key read once a second almost never fires. The rule scales itself: hot keys refresh early, cold keys just expire.
| # | Replay of event 15 | Result, from the reference implementation |
|---|---|---|
| 15d | XFetch, β = 1, Δ = 20 ms, 12,000 reads a second | The first early refresh fires most likely 110 ms before expiry (0.02 × ln(12,000 × 0.02) = 0.02 × ln 240 ≈ 0.110 s), 121 ms on average. Without single-flight, 1.7 more reads also fire during the 20 ms recompute, on average |
| 15d | The same rule for a product read once a second | An early refresh happens before only about 2% of expiries (1 − e^(−0.02)); otherwise the key expires and misses once |
The fixes on equal terms
| Fix | Replica queries per expiry of k7 | Readers who wait | Staleness served | Coordination | For a cold key | For a batch expiring together |
|---|---|---|---|---|---|---|
| None | 231 | Every reader in the first 21 ms, up to 220 ms | None | None | One query per concurrent miss | A herd on every key (event 14) |
| Single-flight per server | 8 | Up to 20 ms | None | None (in-process) | One per server | One per server per key: still one burst |
| Lease | 1 | Up to 25 ms; up to 10 s if the holder dies | None | The cache server | One | One per key: 10,000 keys expiring together still need about 3,800 fills in the first second |
| Stale-while-revalidate | 1 | None | The old value for one refresh (21 ms here) | A lease or single-flight for the refresh | No help: nothing stored to serve | Helps readers, not the fills: they still come together |
| XFetch | About 1 + 1.7 | None | None: it refreshes before expiry | None | No help: rarely read, so rarely early | Spreads the hot keys' refreshes; cold keys still expire together |
| TTL jitter | Not a fix for one key | None | The fix: 83 fills a second instead of 16,200 misses a second at once |
The main line from 20:10:05 runs leases, stale-while-revalidate, jitter on every fill, and fail fast with capped fills (Part 8), by a flag (event 16).
Drill 01, question 1: a restarted cache
The drill asks why a restarted, empty cache melts the database. It's a stampede on every key at once: each key's first readers all miss together. One fill per key (single-flight or a lease) turns that into one query per distinct key requested, which is still far more than the replica takes in the first seconds. Part 8 does that capacity math.
Why didn't k7's herd grow here? The replica was nearly idle at 20:10:00, so the first query found a free connection and landed in 21 ms. If the replica is already busy when a hot key expires (Part 8), the first fill waits in the queue, the window stays open longer, more readers miss, and the herd feeds itself into timeouts and retries.
What to remember from Part 5
- Jitter spreads expiries; one fill per key collapses readers.
- Stale-while-revalidate means nobody waits: serve the old value while one refresh runs.
- Early refresh (XFetch) refreshes hot keys before they expire and leaves cold keys alone.
Part 6. Negative caching: the keys you don't have
A cache normally stores what exists. A key that doesn't exist misses every time, and every miss is a query. The email links to a product that was deleted yesterday.
Trace: a deleted product in an email
k9 was deleted yesterday, and the email still links to it: 500 reads a second (labelled). Side replay 17, on a copy without leases, and the main line (from the reference implementation, one minute):
| # | Design | Replica queries in one minute | What readers get |
|---|---|---|---|
| 17 | No protection (on a copy) | 30,012: every read, 500 a second, 50% of the replica | "No such product", after a 20 ms query each |
| Leases only (the main line after 20:10:05) | 6: one per 10 s token | The holder gets "no such product". The other 30,006 are told to wait for a value that never arrives: the holder found no row and stored nothing | |
| 17a | Negative caching (main line from 20:11:00): store "absent" for k9, TTL 60 s | 1 | "No such product", from the cache |
Negative caching gives the lease's protection without a lease server, and it gives every reader an answer.
A create must replace a cached "absent"
Side row 17b: at 20:11:20 a probe for k12 (not yet created) caches "absent" until 20:12:20. k12 is created at 20:11:30. If the create doesn't invalidate it, readers are told "no such product" until 20:12:20: 504 reads at 10 a second (labelled), from the reference implementation. With Fix 1 it's automatic: the create writes k12 version 1, and "absent" counts as version 0, so it loses. A create is a write like any other.
Random keys never repeat
Side row 17c: a bot asks for 300 random ids a second (30% of the replica, labelled). Nothing repeats (the reference implementation drew 18,000 ids in a minute and saw no repeat), so negative caching stores 300 × 60 = 18,000 useless "absent" entries a minute, which take memory, and still sends 300 queries a second. Enumeration is stopped before the cache: validate the id's format, limit each client (the Rate Limiting loop primitive, coming), or keep a filter of the ids that exist (a Bloom filter) where the id space is known.
What to remember from Part 6
- Cache "absent" when the same missing key is asked for again and again.
- A create must replace a cached "absent", like any other write.
- Random keys never repeat: stop enumeration with validation and limits, not a cache.
Part 7. Sizing and eviction
The main line's cache holds the whole catalog. Then the next budget review halves the cache. What does that do to the database? Side replay M, on a copy: the same 20,000 reads a second, with memory for fewer entries.
Working set and hit ratio
When memory is full, the cache must evict an entry for every new one, and which one it evicts decides the hit ratio. From the reference implementation (eviction only, no TTLs; the sampled policies work as Valkey does, described below):
| Memory holds | LRU (sampled, 3 keys) | LFU (sampled, 3 keys) | Keeping exactly the top N (the best possible) | Replica load under LRU | Under LFU |
|---|---|---|---|---|---|
| 10,000 entries | 74.2% | 78.7% | 81.0% | 20,000 × 0.258 ≈ 5,160/s | 4,270/s |
| 20,000 entries | 81.3% | 84.6% | 86.7% | 3,750/s | 3,080/s |
| 50,000 entries | 91.5% | 92.6% | 94.3% | 1,690/s | 1,490/s |
| 100,000 (all) | 99.2% (only the TTL's misses, Part 1) | 99.2% | 99.2% | 159/s | 159/s |
Exact LRU gives the same as sampled LRU to within half a point (73.8%, 81.2%, 91.8%). Every size below the whole catalog puts the replica over its 1,000 ceiling. Formula 1 is the whole argument: a cache sized for "most" of the working set still sends the database more than it can take.
Synthesizing vector architecture diagram...
What to notice: the curve flattens slowly. Going from half the catalog to all of it moves the hit ratio from 91.5% to 99.2%, and the replica's load from 1,690 to 159 queries a second.
The cache holds 50,000 of 100,000 products under LFU and the hit ratio is 92.6%. Doubling memory costs money every month. How do you decide?
LRU, LFU and scans
- LRU (least recently used) evicts the key read longest ago. Under skew, a cold key that was just read once pushes out a warm key that is read every few seconds.
- LFU (least frequently used) evicts the key read least often. It keeps more of the hot set: 78.7% against 74.2% at 10,000 entries.
Valkey approximates both by sampling. It doesn't keep a list of every key in order. On each eviction it samples a few keys (maxmemory-samples), keeps the best candidates in a small pool, and evicts the best one. ElastiCache's default is 3 samples ("By default, Redis OSS chooses 3 keys and uses the one that was used least recently"); open-source Valkey's valkey.conf uses 5. In the reference implementation 3 or 5 samples changed the hit ratio by less than 1 point. Sampling, and a logarithmic counter that makes many keys look alike, keep the sampled LFU below the best possible 81%.
Scan pollution (side row M2): a nightly job reads all 100,000 products once, in 10 s, through the cache (misses fill it), with memory for 50,000. From the reference implementation, the normal traffic's hit ratio (a separate run with its own random draws, so its "before" figures differ from M1's by a few tenths):
| Policy | Before | During the scan | Next 10 s | The 20 s after |
|---|---|---|---|---|
| LRU (3 samples) | 91.5% | 89.4% | 90.1% | 91.5% |
| LFU (3 samples) | 92.3% | 92.6% | 92.8% | 92.9% |
Under LRU the scan's one-off keys push out part of the hot set, and the hit ratio dips for about 20 s. Under LFU a new key starts with a small counter, so the scan's keys are evicted first. The fix is to keep the job out of the cache: it reads the replica directly, throttled. The job's own misses are replica queries too: about half of 10,000 a second would be 5 times the replica.
Eviction policies
maxmemory-policy | Evicts | Use it for |
|---|---|---|
allkeys-lru / allkeys-lfu | Any key, by approximated LRU / LFU | A pure cache where every key can be refilled |
allkeys-random | Any key, at random | Rarely: when every key is equally likely to be read |
volatile-lru / volatile-lfu / volatile-random | Only keys with a TTL | A cluster mixing refillable keys (with a TTL) and keys that must stay (without one) |
volatile-ttl | Keys with a TTL, the shortest remaining TTL first | When the TTL says what matters least |
noeviction | Nothing: writes fail when memory is full | Data you can't lose (sessions, locks, counters), in its own cluster |
The defaults differ, so check yours: ElastiCache's default parameter groups use volatile-lru; open-source Valkey's valkey.conf uses noeviction.
Side row M3, from the reference implementation: under volatile-lru, a developer's keys written with SET and no EX have no TTL, so nothing is evictable. With memory for 1,000 entries, the 1,001st write and every one after it fails with an out-of-memory error (200 of 1,200 in the run). Keep sessions, locks, counters and dedup keys out of a cache cluster; they need noeviction in a cluster of their own.
Usable memory is less than the node's. ElastiCache reserves memory for background work: reserved-memory-percent defaults to 25%, so on a node with 13.5 GB, "the maxmemory is about 10.1 GB". Tombstones and "absent" entries count against it too.
What to remember from Part 7
- Size the cache to the working set, then check the database's load at the hit ratio you get.
- LFU survives one-off scans; LRU doesn't.
- Know your eviction policy:
volatile-*evicts only keys with a TTL.
Part 8. When the cache is empty: failover, cold starts and warming
The database was sized for a 99% hit ratio. At 20:15:00, in the middle of the sale, shard B's node fails. It holds k7 and 22.7% of the catalog's reads, and it has no replica.
Trace: shard B dies during the sale
| # | Time | Event |
|---|---|---|
| 18 | 20:15:00 | Shard B's node fails: every GET for its keys errors. ElastiCache takes the failed primary offline and provisions a new one; with no replica, "the new primary starts empty and data is lost unless you restore from a snapshot" |
| 18 | 20:17:00 | The replacement node is back, empty (our example's two minutes; AWS gives no figure, only that this is a "longer write interruption" than a failover to a replica) |
The main line already runs the 20:10:05 flag: fail fast, with capped fills (18a below). First, what the naive app would have done.
Cache error is not a miss, on a copy
Side replay 18x: the same failure with the naive rule "cache error = miss: read the database". Every shard-B read goes to the replica: 4,539 catalog reads plus all of k7's 12,000 = about 16,500 reads a second, 16.5 times the replica's ceiling, on top of the other shards' fills. Under the overload model, from the reference implementation:
| Second | Queries sent (with retries) | Answered in time | Answered late | Cancelled unrun | Timeouts | Failed after the retry | Other shards' readers failed |
|---|---|---|---|---|---|---|---|
| 20:15:00 | 16,735 | 983 | 0 | 0 | 0 | 0 | 0 |
| 20:15:01 | 32,713 | 77 | 923 | 14,716 | 15,678 | 0 | 0 |
| 20:15:02 | 34,062 | 0 | 1,000 | 31,692 | 32,713 | 15,678 | 185 |
| 20:15:30 | 38,422 | 0 | 1,000 | 37,114 | 38,118 | 18,892 | 2,539 |
| 20:16:00 | 42,725 | 0 | 1,000 | 41,356 | 42,350 | 21,233 | 4,770 |
| 20:16:59 | 42,970 | 0 | 1,000 | 41,747 | 42,751 | 21,323 | 4,853 |
The other shards' top keys were expiring for the second time between 20:14 and 20:16 (they were refilled 20:04 to 20:06 with a plain 600 s TTL). Once the replica stops answering in time, each of those keys fails to refill and keeps missing, so the load grows past 42,000 queries a second. Over the two minutes: about 4.9 million queries sent; the replica ran 119,983 of them, and 1,060 (0.9% of the queries it ran) came back in time. About 2.4 million readers failed after their retry, on all four shards: a cache outage on one shard became a database outage for the whole shop. Retries doubled the offered load; how retries and shedding should work is the Retries, Timeouts, Backpressure & Load Shedding loop primitive's subject (coming). A busy replica is also where a single hot key's herd grows by itself: its first fill waits in the queue, so more readers miss.
Shard B is back, empty, with k7 on it. With one fill per key, about 1,900 fills are needed in the first second. How long until the replica is below its ceiling, and what would you change before the next time?
Cap the fills
textREAD(k) on app server s main line from 20:10:05 1. GET product:{k} with a short timeout 2. the cache errors (the node is down): a. last_known[s][k] is younger than 5 s -> serve it b. a fill for k is already in flight on s -> wait for that one c. fewer than 2 shard-B fills are in flight on s -> fill from the replica, keep it in last_known, serve d. otherwise -> a degraded page ("price unavailable"); no query 3. never retry the cache in a loop; never send the database an uncapped read connection budget: 8 servers x 2 = 16 of the replica's 20 connections for shard B, leaving 4 (200 queries a second) for the other shards' misses
last_known is each app server's in-process copy of what it read recently, kept 5 s. Event 18a, the main line, from the reference implementation:
| Second | Shard-B catalog reads | From last-known copies | Filled | Degraded | k7 reads | k7 degraded | Replica queries | Queue | Timeouts |
|---|---|---|---|---|---|---|---|---|---|
| 20:15:00 | 4,328 | 2,600 | 615 | 1,100 | 12,415 | 0 | 785 | 0 | 0 |
| 20:15:05 | 4,291 | 1,331 | 684 | 2,245 | 12,098 | 72 | 856 | 1 | 0 |
| 20:15:30 | 4,216 | 1,863 | 666 | 1,665 | 12,287 | 31 | 856 | 6 | 0 |
| 20:16:59 | 4,309 | 1,854 | 688 | 1,764 | 12,011 | 0 | 794 | 1 | 0 |
Over the two minutes, shard B's 515,126 catalog reads were served 44% from last-known copies, 16% by capped fills, and 40% got a degraded page. k7 hardly noticed: every server reads it 1,500 times a second, so its last-known copy is always fresh (99.5% served from it, 0.1% degraded). The fills ran at about 670 a second, not the full 800: a server's slot often sits idle for a few milliseconds between reads that need it. The replica stayed at about 800 queries a second with no timeouts, and the other shards' fills (about 170 a second at 20:15, their top keys' second expiry included) never noticed. The CDN keeps serving /p/k7 from its copy under stale-if-error (Part 10, event 28).
Trace: back, empty
At 20:17:00 the node returns with nothing in it. Every shard-B key misses. With leases (one fill per key) and the same cap, from the reference implementation (event 19; hit ratio of shard B's catalog reads):
| Second | Hit ratio | Fills | Degraded |
|---|---|---|---|
| 20:17:00 | 44.2% | 667 | 1,430 |
| 20:17:01 | 58.8% | 638 | 1,098 |
| 20:17:05 | 72.8% | 553 | 595 |
| 20:17:10 | 81.0% | 493 | 323 |
| 20:17:20 | 87.9% | 368 | 143 |
| 20:17:40 | 93.3% | 241 | 49 |
| 20:18:00 | 96.7% | 130 | 7 |
| 20:18:30 | 98.7% | 57 | 2 |
| 20:19:00 | 99.2% | 33 | 0 |
The hit ratio is back above 99%, and stays there, from 20:18:40: 100 s after the node returned. k7 was filled once in the first milliseconds; its other readers waited on the lease a few milliseconds.
Synthesizing vector architecture diagram...
Snapshot S5, from the reference implementation: the fills that would be needed if every fill succeeded at once. What to notice: the bars stay above the 1,000-a-second line (the replica's ceiling) for the first three seconds, which alone need 4,398 fills, and drop below it from second 3. The capped run in the table above takes them at about 670 a second and serves the rest from last-known copies or a degraded page.
Warming
| # | Replay of event 19 | Result, from the reference implementation |
|---|---|---|
| 19a | The same cap (800 fills a second at most, shared), plus a warm-up job that pre-loads shard B's keys in popularity order from the hot-key list the app logs every minute, using whatever part of the cap the readers leave | 51.4% in the first second, 88.4% at +10 s, 95.7% at +20 s; back above 99% at 20:17:28: 28 s instead of 100 s |
Warming rules: pre-load at a capped rate, most popular first; and when a deploy changes the key format, roll it out gradually or read the old key on a miss, because a new key format is a cold cache for every key.
A replica keeps it warm
Side row 19b: with one replica per shard and Multi-AZ on, ElastiCache "promotes the replica with the least replication lag to primary. Writes can resume within seconds." Shard B keeps its keys; the hit ratio stays near its 99.2%; nothing above happens. The cost is one more node per shard. A cache replica is for warmth, not durability: Valkey replication is asynchronous, and AWS says "when a primary node fails over to a replica, a small amount of data might be lost due to replication lag". For a cache that lost data is not just values. It includes invalidations.
What a failover loses
| # | Time | Event (on the 19b copy, with a replica) |
|---|---|---|
| 21 | 20:14:59.980 | k30003's price changes again, v5 → v6. The tombstone TOMB v6 is acknowledged by shard B's old primary |
| 21 | 20:15:00 | The primary fails. Its last 30 ms of writes (labelled) never reached the replica: TOMB v6 is lost |
| 21 | about 20:15:05 | The replica is promoted. It still holds k30003 v5: stale until its TTL |
| 21 | Fix | After any cache failover, rewind the CDC invalidator by more than the replication lag (60 s, to 20:14:05) and re-apply every tombstone. TOMB v6 lands again; repeats are harmless |
Write-behind data is lost too (event 20, in the main line, with no replica): views:{k7} was last flushed at 20:14:50, so the node failure loses the increments since then, 119,465 views in the reference implementation. That's acceptable only because a trending count can be off.
ElastiCache's durability mode for Valkey 9.0 and later changes this: "synchronous writes ensure zero data loss, while asynchronous writes may lose up to 10 seconds of uncommitted data". That is a durable store's job; the Replication, Quorums & Read-Your-Writes loop primitive covers the trade-off.
What to remember from Part 8
- A cache the database can't live without is part of the database's capacity plan.
- Never turn a cache error into an uncapped database read; cap fills and serve stale.
- A cache replica is for warmth; after any cache failover, replay recent invalidations.
Part 9. Consistency: your own writes, other Regions and DAX
We rewind to 20:00:00 and follow the same price change through the people and places that read it. The first is the person who made it.
Trace: the merchandiser can't see her price
| # | Time | Event |
|---|---|---|
| 22 | 20:00:00.000 | The merchandiser saves $39. The write returns the row's new version, v8 |
| 22 | 20:00:02.000 | She reloads k7's page and sees $49: event 5's stale entry, until 20:10:00.045 |
| 22 | 20:00:02 | She thinks the save failed and clicks Save again. The admin tool compares the form with the row, sees no change and writes nothing: still v8 |
Had the second save written v9, the app's delete would have cleared the stale entry by luck, and the next refill would have read the replica, which has had v8 since .200. Nobody should rely on that.
Read-your-writes through a cache
The writer knows the version its write produced, so it can ask for "at least that version". This is the position token of the Replication, Quorums & Read-Your-Writes loop primitive (Part 4), applied to a cache:
textSAVE(k, fields) returns the committed version, here 8 PREVIEW(k, at_least = 8) 1. e = GET product:{k} 2. e is a value with version >= 8 -> serve it 3. otherwise (older, a tombstone, or empty) -> read the PRIMARY (20 ms) and serve it; a fill, if any, is a compare-and-set the response: Cache-Control: private, no-store under a CDN cache behaviour whose minimum TTL is 0 (Part 10 says why)
The preview costs one primary read per request, which is fine for a merchandiser and far too much for 12,000 viewers.
Trace: the second Region
eu has its own app servers, its own Valkey, and a cross-Region replica 1 s behind. The home Region broadcasts each invalidation to eu, where it arrives in 80 ms.
Synthesizing vector architecture diagram...
What to notice: the invalidation reached eu 920 ms before the data did, so the first refill after it read the old price. The fix makes the invalidation wait for the data: eu invalidates from its own replica's stream.
| # | Replay | Result, from the reference implementation |
|---|---|---|
| 23 | Broadcast DEL, no versions | eu's fill at .321 stores $49 v7, stale until 20:10:00.321 |
| 23 | Broadcast TOMB v8 (Fix 1) | The fill is refused (7 < 8). The reader serves its value once, marked stale, or reads the home Region |
| 23 | eu invalidates from its own replica's stream | The delete comes at 20:00:01.050, after eu's replica has v8. Reads before it hit the old entry (staleness = eu's lag + the invalidator's lag, about 1.05 s); the first miss after it reads v8 |
The memcache paper ran one invalidation daemon per database for the same reason ("mcsqueal"), and for a writer's own reads it added a remote marker: the writer sets a marker in its Region before writing to the master Region; a later miss that finds the marker reads the master Region instead of the local, possibly stale replica. ElastiCache Global Datastore replicates a cache across Regions asynchronously, so the secondary Region's copy also receives every invalidation after the replication lag: in the news feed loop's words (step 3.2), a copy of a cache is a second cache.
DAX is in front of your table and you updated an item through it. Which reads can still return the old price, and for how long?
DAX
Side replay D1, the same product on DynamoDB behind DAX, both TTLs at their default of 5 minutes:
| Time | Request | Answer |
|---|---|---|
| 19:58:40 | Query "sale items" through DAX | Cached in the query cache until 20:03:40 |
| 20:00:00 | UpdateItem k7 $39 through DAX | Written to DynamoDB, then DAX "updates the items in its item cache, regardless of the TTL value" |
| 20:00:01 | GetItem k7 through DAX | $39, from the item cache |
| 20:00:01 | Query "sale items" through DAX | $49 until 20:03:40: writes never touch the query cache |
| 20:00:30 | The pricing service writes k5 directly to DynamoDB | DAX's item cache keeps the old k5 until its item TTL. A GetItem with ConsistentRead goes to the table and sees the new one |
More DAX facts that decide designs: a TTL of 0 on the item cache means entries are refreshed only by eviction or write-through, and on the query cache means nothing is cached; DAX caches negative results in both caches (an item that wasn't found stays "not found" until its TTL, or until it's written through DAX); TransactGetItems passes through like a strongly consistent read; a hit never reaches the table, so it uses no table read capacity, and a miss is an ordinary eventually consistent read.
What to remember from Part 9
- Read-your-writes through a cache: carry the version your write returned and skip older entries.
- In another Region, invalidate after that Region's replica has the write.
- DAX's item cache follows writes made through DAX; its query cache doesn't follow writes at all.
Part 10. The edge: CDN, browser and HTTP caching
We rewind to 20:00 once more. Valkey is only one of the copies: the CDN's edges, its shield and every browser keep their own, each with its own clock. And a CDN that doesn't recognize two requests as the same object doesn't cache at all.
Back to 20:00, on a copy: every link is a different object
Side replay 24, on a copy where the email's links carried ?utm_campaign=flash&sub=<subscriber id> and the cache policy kept every query string in the cache key. (The main line's policy keys on the path only.) The 12,000 page views a second are requests that reach the CDN; a reload that the viewer's browser answers from its own copy isn't one. From the reference implementation, one minute:
| What | Result |
|---|---|
Edge hit ratio for /p/k7 | About 0%: each link is its own object, and a subscriber who comes back within 30 s is answered by their own browser, never by the edge |
| Requests reaching the origin | about 12,000 a second |
| App tier | 20,000 + 12,000 + 12,000 = 44,000 a second against 40,000: 110% |
product:{k7} on shard B | 12,000 private reads + 12,000 renders = 24,000 a second on one key: the per-key ceiling of the Sharding, Hot Keys & Rebalancing loop primitive |
| # | Replay | Result, from the reference implementation |
|---|---|---|
| 24a | The main line's cache policy: path only, no query strings, no cookies | 24 origin fetches in 120 s: 0.2 a second, one per 5 s at the shield, whatever the viewers do |
Synthesizing vector architecture diagram...
What to notice: in the main line, identical requests are collapsed twice (at each edge and at the shield), so the origin sees one fetch every 5 s. In the "On a copy" panel, the key includes the subscriber id, so nothing is identical and nearly every view reaches the origin.
CloudFront collapses simultaneous requests: "if there are simultaneous requests for the same object … CloudFront pauses before forwarding the additional requests", and it "only collapses requests that share a cache key". It even collapses under no-cache; only a minimum TTL of 0 together with private, no-store, no-cache, max-age=0 or s-maxage=0 prevents it. Origin Shield adds one more layer in front of the origin, consolidating misses from all edges to "as few as one request going to your origin".
Cache-Controldirectives, and the Ageheader
| Directive | Who obeys it | Meaning |
|---|---|---|
max-age=N | Browsers (and shared caches, when there's no s-maxage) | Fresh for N seconds |
s-maxage=N | Shared caches only (CloudFront, the shield) | Fresh for N seconds; CloudFront uses it, clamped between the cache policy's minimum and maximum TTL |
private / no-store | private: only the browser may store it. no-store: nobody may | For personal data. On CloudFront, only with a minimum TTL of 0 (below) |
stale-while-revalidate=N | Caches that support it (CloudFront does) | After it turns stale, serve it for up to N seconds more while one refresh runs |
stale-if-error=N | Caches that support it (CloudFront does) | If the origin "is unreachable or returns an error code that is between 500 and 600", serve the stale copy for up to N seconds |
Age: N (a response header) | Set by caches | How old the copy already is; the next cache counts its freshness from there |
stale-while-revalidate and stale-if-error are defined in RFC 5861; the others, with Age and Vary, in RFC 9111 (HTTP Caching). CloudFront caps both stale windows at the cache policy's maximum TTL. Our origin sends public, max-age=30, s-maxage=5, stale-while-revalidate=10, stale-if-error=300.
The cache key
Put in the cache key everything the response depends on, and nothing else. On CloudFront the key is what the cache policy names: the path, plus the query strings, headers and cookies it lists. Every value you add splits the objects and lowers the hit ratio, as event 24 showed. Every value you leave out that the response depends on serves one viewer's answer to another (side row 25a, in Part 13: the page renders euros or dollars by Accept-Language, an origin request policy forwards that header to the origin but leaves it out of the cache key, so each edge stores whichever currency filled it first). CloudFront doesn't build the key from the origin's Vary header; it removes some Vary values before replying, and without a policy that forwards it, it removes Accept-Language before the origin ever sees it. Cookies you forward become part of the key, and CloudFront caches the origin's Set-Cookie too.
The app cache became right at T. The shield and the edges each cache for 5 s with stale-while-revalidate 10 s, and browsers for 30 s. When is the last moment anyone sees $49, with and without Age passed on, and what one header change would cut it most?
Snapshot S6, the last moment each layer served $49 in the main line, from the reference implementation (worst case over request timings, high traffic):
| Layer | With Age passed on | Without Age |
|---|---|---|
| Valkey (event 5's entry) | 20:10:00.045 | 20:10:00.045 |
| Origin Shield | 20:10:05.0 | 20:10:10.1 |
| Edges E1 and E2 | 20:10:05.1 | 20:10:15.1 |
| Browsers | 20:10:30.0 | 20:10:45.1 (20:11:00 at low traffic) |
A browser that honors stale-while-revalidate can show the old copy once more, up to 10 s later, while it refreshes in the background. CloudFront sends Age on cache hits: we saw it on live responses (age: 692387 on a Hit from cloudfront response). AWS doesn't document how its tiers age a copy fetched from another tier, so plan with the column without Age. With Fix 1 in place from the start, Valkey would have been wrong for milliseconds, not ten minutes, and each tier's window would start at 20:00:00.
A personal page cached publicly
| # | Time | Event |
|---|---|---|
| 25 | 20:02:00 | A developer adds "Hi Alice, your member price is $35" to /p/k7's HTML, still public, s-maxage=5. Alice's request is the first after the deploy, so her page fills the shield and both edges. From the reference implementation, about 30,000 views at E1 and 30,000 at E2 see Alice's name and price, about 5 s of each edge's traffic, until the next refresh brings someone else's |
The fix has two parts. Personal data goes only in the private request /api/me/k7, which answers Cache-Control: private. And the CDN cache behaviour for /api/me/* must have a minimum TTL of 0: "If your minimum TTL is greater than 0, CloudFront uses the cache policy's minimum TTL, even if the Cache-Control: no-cache, no-store, and/or private directives are present in the origin headers." The managed CachingOptimized policy has a minimum TTL of 1 s, so it would cache a private response for a second.
Purge or version
At 20:20:00 k7's photo changes (event 27). The photo gets a new, versioned URL, k7.9f3a.jpg, cached for a year; the page now points at the new name, and nothing is invalidated. AWS recommends this: "Versioning is less expensive". Now the arithmetic if you purged on every price change instead:
| What | Arithmetic | Result |
|---|---|---|
| Paths a month | 50,000 price changes a day × 30 days | 1,500,000 |
| Price | The first 1,000 paths a month are free, then 0.005 | $7,495 a month |
| Rate | CloudFront accepts 150 paths or tags a second, and one wildcard invalidation a second. AWS doesn't say whether that's per distribution or per account, so budget it as shared by every purge you send | A burst of 50,000 takes 50,000 ÷ 150 ≈ 333 s just to be accepted |
| Reach | CloudFront "forwards the request to all edge locations within a few seconds" | Browsers keep their copies: a purge never reaches them |
A wildcard path (/p/*) counts as one path, and cache-tag invalidation (one tag per product) also counts one path per tag. Use versioned URLs for files, short TTLs for data, and purges for takedowns. That answers the drill The Image Service That Made Every Origin Server Sweat (question 1): with a versioned URL, the new photo shows as soon as a viewer gets a page that points at it (at most the page's own TTLs); with a purge, the edges drop the old one within seconds, but browsers keep it for its max-age. Takedowns that must reach further, like the YouTube loop's harmful content (step 3.4), add an edge blocklist and invalidation by cache tag; ISP caches outside your CDN are beyond any purge.
Serving stale at the edge
During shard B's outage (event 28, 20:15:00 to 20:17:00), the origin's render of /p/k7 doesn't use the last-known copy: on a cache error it answers 503, so the edge's copy is the fallback, rather than a degraded page the edge would then cache for 5 s. The edge keeps its copy under stale-if-error=300. From the reference implementation: no viewer got an error, and the oldest copy served was 125 s old, inside 5 + 300. CloudFront also serves an expired copy when the origin is unreachable and the minimum or maximum TTL is above 0, even without the directive; and it caches an error response for 10 s by default, so an early 5xx or 404 from a half-deployed origin doesn't stick for long.
A stale copy doesn't always need a new body: if it carries an ETag, the cache can revalidate it with If-None-Match, and the origin answers 304 Not Modified with no body.
Why an edge at all
The drill's second question: why a CDN rather than a bigger origin? A bigger origin can take more requests, but it can't remove the distance: every viewer far from it pays the round trips of TLS setup and the request itself. An edge answers from nearby and collapses identical misses, so the origin's load stops depending on the number of viewers (24a: 0.2 fetches a second for 12,000 views). The costs: every request is billed, hit or miss, so a CDN saves origin capacity, egress and latency, not request fees; and it adds the staleness above.
What to remember from Part 10
- Put in the cache key everything the response depends on, and nothing else.
- Staleness adds across tiers unless each passes on its age.
- Use versioned URLs for files and short TTLs for data; purge only for takedowns.
Part 11. End to end through the layers
Each layer caches something, has its own staleness bound, learns about a change in its own way, and fails in its own way. Follow one page view of k7 at 20:30:00 from a phone, and then the next price change, k7 back to $49 as v9 at 21:00:00, with every fix on.
Synthesizing vector architecture diagram...
What to notice: the public page and the private request take different paths from the phone. The public one stops at the first cache that has it; the private one always reaches the app and Valkey. The dotted line is the only path to the database, and only one fill per key takes it.
Synthesizing vector architecture diagram...
What to notice: two paths write the same tombstone at home, and eu waits for its own replica before invalidating. The HTTP tiers are never told: their TTLs are the only bound, and the browsers' is the longest.
The price change, from the reference implementation (event 29):
| Time | Event |
|---|---|
| 21:00:00.000 | The primary commits k7 $49 v9; the save returns v9 to the admin tool |
| 21:00:00.005 | The app writes TOMB v9 to shard B (the fast path) |
| 21:00:00.009 | The first miss gets the lease; its fill reads the replica (still v8): refused, 8 < 9 |
| 21:00:00.050 | The CDC invalidator writes TOMB v9 again: no change (a repeat) |
| 21:00:00.051 | The lease holder, re-reading the primary, stores $49 v9; the readers who waited on the lease get it |
| 21:00:00.200 | The home replica applies v9 |
| 21:00:01.000 / 01.050 | eu's replica applies v9; eu's invalidator writes TOMB v9 in eu's Valkey |
| by 21:00:05 | Shield and edge copies turn stale (s-maxage 5 s) and are refreshed by the next request |
| by 21:00:30 | Browsers' max-age runs out (with Age passed on; up to about 21:00:45 without) |
| Hop | What it caches | Staleness bound | What invalidates it | How it fails | What protects the database |
|---|---|---|---|---|---|
| Browser | The public page | max-age 30 s | Nothing can | Shows the old page until max-age | Nothing needed |
| Edge E1 | The public page, per cache key | s-maxage 5 s (+10 s stale-while-revalidate) | A purge (takedowns only) | Serves stale under stale-if-error when the origin fails | Collapsing identical misses |
| Origin Shield | The same, for all edges | 5 s (+10 s) | A purge | Same | One fetch per key per 5 s |
| Load balancer | Nothing | Sends requests to healthy servers | |||
| App server | A last-known copy of what it read, 5 s | 5 s | Its TTL (or a broadcast, if you run one) | Restart: empty | Serves during a cache outage (Part 8) |
| Valkey shard B | product:{k} with its version | TTL 600 s ± 10% (soft), +60 s (hard) | Tombstones from the app and the CDC invalidator | Node loss: empty, or a replica promoted with a lost tail | One fill per key (lease), stale-while-revalidate, capped fills |
| Read replica | The whole table, a log behind | Its lag (200 ms here) | It follows the primary's log | Lag grows; refills may read old data (race B) | Its 1,000 queries a second is the budget everything above protects |
| Primary | The truth | None | Only writes, checkout, previews and refused fills reach it |
The app server's in-process copy is the same idea as the Sharding, Hot Keys & Rebalancing loop primitive's micro-cache (Part 4). If in-process copies must follow changes faster than their TTL, Valkey's server-assisted client-side caching (the CLIENT TRACKING command) sends invalidation messages to the clients that hold copies.
What to remember from Part 11
- Every layer is a copy with its own clock and its own way to learn of a change.
- The database is protected by one fill per key and a cap on fills, not by hope.
- Write down the staleness bound of every layer; the user sees their sum.
Part 12. On AWS
Every AWS cache has a TTL, an invalidation story, and a failure mode. Stated only from AWS's public documentation.
Managed services that use it
| Service | What it caches, and where | Staleness and invalidation | Failure behaviour | What AWS states |
|---|---|---|---|---|
| Amazon ElastiCache (Valkey, Redis OSS, Memcached; node-based or Serverless) | The application cache; cache-aside is your code | A TTL per key; invalidation is your app's or your invalidator's; eviction by maxmemory-policy, volatile-lru in the default parameter groups; maxmemory-samples 3; reserved-memory-percent 25% | Multi-AZ: "promotes the replica with the least replication lag", "writes can resume within seconds", "a small amount of data might be lost"; without replicas "the new primary starts empty". Node-based Memcached: "a node failure will always result in some data loss". Serverless: a replicated Multi-AZ architecture. Valkey 9.0+ durability mode: synchronous writes lose nothing, asynchronous "may lose up to 10 seconds". Global Datastore replicates across Regions asynchronously | Its caching-strategies page: with lazy loading "each cache miss results in three trips", write-through data "is never stale" (one writer), and TTLs keep data from getting too stale |
| DynamoDB Accelerator (DAX) | An item cache (GetItem, BatchGetItem) and a query cache (Query, Scan) in front of DynamoDB | Item cache: write-through for writes made through DAX; query cache: not updated or invalidated by writes; both default to a 5-minute TTL; negative results cached in both; LRU eviction | A replica takes over if the primary node fails; nodes replicate eventually consistently | Strongly consistent reads and TransactGetItems pass through and aren't cached; after TransactWriteItems, DAX fills its item cache with a background TransactGetItems; writes that bypass DAX (including global tables' replicated writes, page 06) aren't seen |
| Amazon CloudFront, with Origin Shield | Edge locations ("750+ POPs"), 15 regional edge caches, and an optional Origin Shield | The cache policy's minimum, default (86,400 s) and maximum (31,536,000 s) TTLs clamp s-maxage and max-age; stale-while-revalidate and stale-if-error supported, capped by the maximum TTL. Invalidation: 1,000 paths a month free, then $0.005 a path; 150 paths or tags a second and one wildcard invalidation a second; "within a few seconds" to all edge locations | Serves stale content under stale-if-error for 5xx or an unreachable origin; also serves the previous copy when the origin is unreachable and the minimum or maximum TTL is above 0; caches error responses for 10 s by default | The cache key is the cache policy's query strings, headers and cookies; a minimum TTL above 0 overrides private and no-store; Vary doesn't build the key; request collapsing per cache key; Origin Shield consolidates to "as few as one request"; versioned file names are "less expensive" than invalidation |
| Amazon API Gateway (REST APIs) | A per-stage response cache you size | Default TTL 300 s, maximum 3,600 s, 0 disables; only GET cached by default; a client can invalidate with Cache-Control: max-age=0 if allowed by execute-api:InvalidateCache, and "if you don't impose an InvalidateCache policy … any client can invalidate the API cache" | REST APIs only: HTTP APIs have no cache; charged by the hour and "not eligible for the AWS Free Tier" |
Running it yourself
| Option | What it is | Facts and sizing |
|---|---|---|
| Amazon EC2 running Valkey, Redis or memcached (or on EKS), and Varnish or nginx as an HTTP cache | You own the eviction policy, replicas, failover, warm-up, the invalidator and the fill script | Valkey's own defaults differ from ElastiCache's: maxmemory-policy noeviction, maxmemory-samples 5, and "LRU, LFU and volatile-ttl are implemented using approximated randomized algorithms" (valkey.conf) |
| Sizing in words | Memory = working set × entry size, plus overhead and reserved memory; replicas per shard for warmth; fills capped so the database survives a cold shard; nodes in several AZs, sized to lose one | A client in another AZ pays $0.01 per GB on the EC2 side for traffic to and from ElastiCache. Measure latency; don't assume it |
Look-alikes that are not this mechanism
| Look-alike | Why it looks like a cache | Why it isn't |
|---|---|---|
| Amazon MemoryDB | Valkey and Redis OSS compatible, in memory | A durable database: "Successful write operations are durably stored in a distributed Multi-AZ transactional logs before returning to clients", positioned as a primary database "eliminating the need to separately manage both a cache and durable database". As your source of truth there's nothing to invalidate; used as a cache, you pay for durability you don't need |
| RDS and Aurora read replicas | "A read copy" | They follow the primary's log: a replica has a subscription, a cache doesn't. Their staleness is lag, not a TTL (page 06) |
| DynamoDB itself, and DynamoDB TTL | "Fast enough that we don't need a cache"; "TTL expires items" | Per-partition limits still apply (page 08). TTL deletes expired items "within a few days of their expiration time": a cleanup, not an expiry a reader can trust. Check expires_at in code |
| The database's buffer cache (Aurora, RDS) | "It caches pages in memory" | Inside the engine and always consistent with it; nothing to invalidate |
| ElastiCache used as a store (sessions, rate-limit counters, locks, dedup keys) | The same service | Not refillable from a source: it needs noeviction and a cluster of its own (Part 7) |
| AWS Global Accelerator, S3 Transfer Acceleration | "An edge network, faster reads" | They route traffic over AWS's network; they cache nothing |
DNS answers are caches (a TTL per answer, with the same stacking rule), and they matter most when traffic moves between Regions: the Multi-Region Failover loop primitive covers the DNS tail (Part 5).
What to remember from Part 12
- ElastiCache, DAX, CloudFront and API Gateway are caches: each has a TTL, an invalidation story and a failure mode.
- MemoryDB is a durable database with a cache's API.
- Managed or not, you still own versions, invalidation and the empty-cache day.
Part 13. What you've learned
Back to the kettle
At 20:00:00 the price dropped to 49 for ten minutes, to about 6.8 million page views. And the database, sized for a 99% hit ratio, would take 32 times its ceiling on the day the cache was empty. The pieces:
- A cache is a copy with no subscription (Part 1). Its load on the database is reads × (1 − hit ratio), and each key read costs at most one miss per TTL.
- Where the cache sits decides what can go wrong (Part 2): write-through still crosses with two writers, write-behind makes the cache the truth until it flushes.
- The ten minutes of $49 came from a refill that read a lagging replica and landed after the delete (Part 3). A tombstone with a version and a compare-and-set fill refused it; a lease would have stopped only a slow reader.
- Invalidations get lost (Part 4): the change stream retries them, and invalidating from the copy you fill from orders them.
- Herds: jitter kept 10,000 expiries at 225 queries a second instead of a collapse where 80 answers in a minute arrived in time; one fill per key took
k7's expiry from 231 queries to 1 (Part 5). Negative caching gave the deleted product one query a minute (Part 6). - Memory: half the catalog in memory still put the replica at 1,490 to 1,690 queries a second (Part 7).
- The empty day: "cache error = miss" answered 0.9% of the queries the replica ran in time; fail fast with capped fills kept the replica under its ceiling, the cold shard was back to 99% in 100 s, 28 s with warm-up, and a replica per shard avoids the cold start (Part 8).
- Readers who must see the write carry its version; another Region invalidates after its own replica has the data; DAX's query cache follows no writes (Part 9).
- The edge: a cache key with a subscriber id cached almost nothing; the HTTP tiers added 30 to 60 s after the app was right; purging every price change would cost $7,495 a month (Part 10).
- The costs: a version on every row and a script per fill, an invalidator and its lag, stated TTLs at every tier, a capacity plan for the empty day, and one more node per shard.
The whole story, event by event
From the reference implementation:
| # | Time | Event | Result |
|---|---|---|---|
| 1 | 19:59:00 | Baseline, TTL 600 s | At most 167 misses/s; 159 in steady state (99.2%); 142/s at 19:59 after the pre-warm (14% of the replica) |
| 1a | side | The hook at 32,000 reads/s | 99% → 320, 98% → 640, 97% → 960, 0% → 32,000 queries/s |
| 2 | 19:59:10 | A read of k7 | Hit on shard B: $49 v7, until 20:04:07.0 |
| 13 | 19:55:00 | Pre-warm, TTL 600 s ± 10% | Expiries 20:04:00 to 20:06:00: 83 extra fills/s, 225 queries/s in all, no timeouts |
| P1 | side | Six designs, one price change | Part 2's table |
| P2 | side | Write-through, two writers | Cache 38 v9, until 20:10:00.018 |
| P3 | side | Write-behind for the price | The database says $49 until the flush; a node loss loses the change |
| 3 | 20:00:00.000 | Commit $39 v8 | The replica applies it at .200 |
| 4 | 20:00:00.005 | DEL product:{k7} | Empty |
| 5 | .011 to .045 | Four misses read the replica | $49 v7 stored; t₅ = 20:00:00.045, until 20:10:00.045 |
| 6 | to 20:10:00.045 | Stale | About 6.8 million page views see 39 |
| 6a | copy | Race A: a reader of the primary pauses | $49 v7 stored at .050, after the delete |
| 6b | side | Refill conditions | DEL + NX: stale; overwrite + plain SET: stale (race A); overwrite + NX: closed |
| 6c | side | Delete before commit | A fill that read v7 lands at .018: stale |
| 6d | side | Delayed double delete (+500 ms) | Clears event 5's fills; not a 1 s pause (its SET lands at 1.010) |
| 7 | replay | Fix 1: TOMB v8 + compare-and-set | A1 refused at .032, stores $39 v8 at .053; 12 refusals, 24 queries |
| 7a | side | Leases | Race A refused; race B stored |
| 7b | side | No TTL | Event 5's entry stays until an eviction or the next change |
| 7c | side | Two writers, writer 1's TOMB v8 delayed past TOMB v9 | Unconditional tombstone write: 38 v9; compare-and-set write: v8 refused, $38 v9 stored |
| 22 | 20:00:02 | The merchandiser reloads | Sees $49; her second save writes nothing; fix: read at least v8 |
| 23 | .080 to 1.050 | The eu Region | Broadcast DEL: stale until 20:10:00.321; TOMB v8: refused; eu's own stream: 1.05 s |
| D1 | side | DAX | Item cache 49 until 20:03:40 |
| 24 | copy | ?utm…&sub=<id> in the cache key | Edge hit ratio about 0%; about 12,000 origin requests/s; app tier 110%; 24,000 reads/s of product:{k7} |
| 24a | main line | Path-only cache key | 0.2 origin fetches/s |
| — | 20:02:00 | Flag: Fix 1 on | |
| 25 | 20:02:00 | A personal greeting in the public page | About 30,000 views at each edge see Alice's name and price |
| 25a | side | Accept-Language not in the key | Each edge stores the first filler's currency |
| 8 | 20:03:00.000 | k30003 v5; the app's tombstone times out | 35 stale reads; fixed at 20:13:06.797 (607 s) |
| 8a | replay | CDC invalidator | TOMB v5 at .051; the next read (20:03:17.822) fills v5 |
| 8b | side | Invalidate from the replica's stream | Delete at .250 closes race B; race A (a SET at .260) stays open |
| 8c | copy | Bulk import, 50,000 prices in 20 s | Invalidator lag peaks at 30 s |
| 8d | side | A new key per version | Old page for at most the pointer's 1 s; one cold fill per bump |
| — | 20:04:00 | Flag: CDC invalidator on | |
| 14 | copy | Pre-warm without jitter | 8,243 misses in the first second; from the third second no answer in time; 80 of 59,000 in a minute; 644 keys refilled; no recovery |
| 15 | 20:10:00.045 | k7 expires under 12,000 reads/s | 231 queries; queue 202; empty after 270 ms; p99 220 ms; no timeouts; $39 v8 stored |
| 15a | replay | Single-flight per server | 8 queries |
| 15b | replay | Lease | 1 query; readers waited up to 25 ms |
| 15c | replay | Stale-while-revalidate | 1 query; nobody waits; 231 reads get $49 during the refresh |
| 15d | replay | XFetch | First early refresh 110 ms before expiry (most likely), 121 ms on average; 1.7 duplicates; about 2% for a key read once a second |
| 26 | 20:10:00.045 on | TTL stacking | Last $49 at +30 s with Age; +45 s (high traffic) or +60 s (low) without |
| 16 | 20:10:05 | Flag: leases, stale-while-revalidate, jitter on every fill, fail fast with capped fills | |
| 17 | copy | Deleted k9, 500 reads/s, no leases | 30,012 queries a minute (50% of the replica) |
| 17a | 20:11:00 | Negative caching | 1 query a minute (leases alone: 6, with readers left waiting) |
| 17b | 20:11:20 | k12 "absent", created 20:11:30 | 504 reads told "no such product" unless the create invalidates |
| 17c | side | 300 random ids a second | 18,000 queries and 18,000 useless entries a minute; no repeats |
| M1 | copy | Memory for 10,000 / 20,000 / 50,000 | LRU 74.2 / 81.3 / 91.5%; LFU 78.7 / 84.6 / 92.6%; best possible 81.0 / 86.7 / 94.3% |
| M2 | copy | A nightly scan through the cache | LRU dips from 91.5% to 89.4%; LFU doesn't |
| M3 | copy | volatile-lru, keys without TTL | Writes fail once memory is full |
| 18 | 20:15:00 | Shard B's node fails, no replica | Back at 20:17:00, empty |
| 18x | copy | "Cache error = miss" | About 16,500 reads/s from shard B, growing past 42,000 queries/s as the other shards' keys fail to refill; 0.9% of the queries the replica ran answered in time; about 2.4 million readers failed |
| 18a | 20:15:00 to 20:17:00 | Fail fast, capped fills | Shard B: 44% last-known, 16% filled, 40% degraded; k7 99.5% last-known; replica about 800/s, no timeouts |
| 20 | 20:15:00 | Write-behind counter lost | 119,465 view increments |
| 28 | 20:15:00 to 20:17:00 | stale-if-error at the edge | No viewer errors; oldest copy served 125 s |
| 19 | 20:17:00 | Back, empty | 1,901 fills needed in the first second, below 1,000 from the fourth; capped: 99% at 20:18:40 (100 s) |
| 19a | replay | Warm-up from the hot-key list | 99% at 20:17:28 (28 s) |
| 19b | copy | A replica per shard | Promoted within seconds, warm |
| 21 | copy (19b) | Lost tombstone for k30003 v6 | Stale on the promoted replica; rewind the invalidator 60 s |
| 27 | 20:20:00 | k7's photo | Versioned URL; purging every price change: $7,495 a month |
| 29 | 20:30:00 / 21:00:00 | One page view; k7 $49 v9 | $49 v9 stored at .051; eu at 1.050; browsers by about +30 s |
The cheat card
| Topic | Remember |
|---|---|
| What a cache is | A copy with no subscription: fresh by expiry or by invalidation |
| Load | Database queries = reads × (1 − hit ratio); size the database for the hit ratio you can lose |
| The TTL's miss cost | At most one miss per key read per TTL; a hot key almost exactly one |
| Patterns | Cache-aside by default; write-through still needs versions; write-behind only for data you can lose |
| Invalidation | After the commit; leave something the refill can check: the new value (refill SET NX) or a versioned tombstone; with versions, both the invalidation and the fill are compare-and-sets |
SET NX | First writer wins, whatever its age: useless after a bare DEL |
| Leases | One token per key per 10 s; a delete voids tokens; closes the slow reader, not the lagging replica |
| Sources | App fast path + change stream guarantee + TTL bound; invalidate after the copy you fill from |
| Stampede fixes | Jitter for batches; single-flight or a lease for one key; stale-while-revalidate so nobody waits; XFetch refreshes hot keys early |
| Negative caching | For repeated misses only; a create must replace "absent" |
| Eviction | volatile-* evicts only keys with a TTL; ElastiCache default volatile-lru, Valkey default noeviction; LFU survives scans |
| Empty cache | Fail fast, serve last-known or degraded, cap fills per server; pre-load at the cap; replicas for warmth; replay invalidations after a failover |
| Read-your-writes | Carry the write's version; skip older entries; read the primary |
| Regions | Invalidate from each Region's own replica stream, or compare versions |
| DAX | Item cache write-through (via DAX only); query cache ignores writes; 5-minute TTLs; strong reads pass through |
| HTTP | s-maxage for shared caches, max-age for browsers; staleness adds across tiers unless Age is passed on; min TTL above 0 overrides private |
| Cache key | Everything the response depends on, nothing else |
| Purge | $0.005 a path after 1,000 a month; 150 a second; never reaches browsers; version file URLs instead |
Failure checklist
- Does every copy (browser, edge, shield, app, cache, replica) have a stated staleness bound, and does someone know their sum?
- Does every invalidation run after the commit, and leave a new value or a versioned tombstone (written with a compare, so an older one can't replace a newer one) rather than a bare delete?
- Is every refill conditional (
SET NXagainst an overwrite, or a version compare), never a plainSET? - Is the change stream the guarantee behind the app's invalidation, with an alarm on the invalidator's lag?
- Do fills and invalidations read the same copy, or does the invalidator wait for the replica the fills read?
- Does every hot key have one fill per key (single-flight or a lease), and do batch loads use TTL jitter?
- Is the database sized for the day the hit ratio drops, with fills capped per server and a degraded answer ready?
- Does a cache error fail fast instead of becoming an uncapped database read with retries?
- Do cache shards the database can't live without have replicas, and is the invalidator rewound after every cache failover?
- Is the eviction policy what you think it is, and are sessions, locks and counters in a cluster of their own?
- Does the CDN's cache key hold everything the response depends on and nothing else (no tracking parameters, no per-viewer values)?
- Is personal data only in
privateresponses, under a CDN behaviour with a minimum TTL of 0?
Think-first drills
Drill 1. A product gets 5,000 reads a second, a fill takes 40 ms, and there are 10 app servers. How many queries does one expiry send with no protection, with single-flight per server, and with a lease? With XFetch at β = 1, about how long before expiry does the first early refresh fire?
Drill 2. A cache holds 1 million keys with a 1-hour TTL, each read at least once an hour; reads are 50,000 a second. What's the worst hit ratio the TTL allows, and the most load it puts on the database? Then halve the TTL.
Drill 3. A price change is committed at t = 0. The app deletes the cache entry at 5 ms; fills read a replica 300 ms behind; a garbage-collection pause can hold a reader up to 1 s. Which of these closes the race: a second delete at 500 ms, a lease, a versioned tombstone, reading the primary?
Interview questions
| Question | Model answer |
|---|---|
| Walk me through cache-aside, and the race that can leave a stale value after a write. How do you close it? | Reads check the cache, and on a miss read the database and fill with a TTL; writes commit, then invalidate. The race: a fill that read the old value (from a lagging replica after the delete, or from the primary before the commit, then paused) lands after the delete and stays until the TTL. A bare delete leaves nothing to check against. Close it with versions on both sides: the invalidation writes a tombstone with the new version and every fill is a compare-and-set that refuses older versions; or, with a single writer, overwrite with the new value and refill with SET NX. Leases stop the slow reader but not the lagging replica. Keep the TTL as the last bound. |
| A hot key expires under 20,000 reads a second. What happens, and what are your options? | Every read before the first fill lands misses: about 20,000 × the fill time queries, a burst that queues on the database; if it's already busy, the window grows into timeouts and retries. Options: one fill per key (single-flight per process, or a lease across the cluster), stale-while-revalidate so readers get the old value while one refresh runs, or probabilistic early refresh (XFetch) so a hot key refreshes before it expires. Jitter is for many keys expiring together, not one. A longer TTL only moves the stampede. |
| Your cache cluster restarts empty at peak. What happens to the database, and how do you design so that isn't an outage? | The database gets the full read rate, far over its ceiling, and if the app treats cache errors as misses with retries, almost no answer arrives before its reader gives up. Design: fail fast on cache errors, serve a last-known copy or a degraded page, cap fills per server so the database stays under its ceiling, use one fill per key, pre-load the hottest keys at the cap, and give shards replicas so a node loss is a failover to a warm copy. After any cache failover, replay recent invalidations from the change stream. |
| When would you use write-through or write-behind instead of cache-aside? | Write-through when readers need the new value right after a write and you accept caching every written key; it still needs version compares because two writers' cache writes can cross, and a failed cache write leaves it behind. Write-behind only for data you can lose or rebuild, like a view counter: until the flush the cache is the only copy, and a node loss loses the writes. Cache-aside stays the default: the database is the truth and a cache failure only costs misses. |
| How do you invalidate caches in two Regions after one write? | Don't broadcast a delete at commit time: it can arrive before the other Region's replica has the write, and the next refill reads the old value. Invalidate in each Region from its own replica's change stream, after the data has arrived, and compare versions on every fill so an older value is refused anyway. For the writer's own reads, carry the version the write returned, or use a marker that sends its reads to the home Region. |
| Your CDN hit ratio fell to near zero after a marketing email. Why, and what do you check first? | The cache key. The email's links probably carry tracking parameters (utm_*, a subscriber id) and the cache policy includes query strings, so every link is a different object, and the origin gets nearly every view. Check the cache policy's query strings, headers and cookies; key on the path and only the parameters the response depends on, or strip tracking parameters at the edge. Then check that no per-viewer value (a cookie, a header) entered the key, and that personal data isn't in the public page. |
Where to go next
- Change Streams & the Transactional Outbox: the stream that feeds the invalidator, order per key, and why caches store "gone" (Parts 3, 4 and 8).
- Replication, Quorums & Read-Your-Writes: the replica's lag behind race B, position tokens (Part 4), and why asynchronous copies lose their tail.
- Sharding, Hot Keys & Rebalancing: one hot key, the micro-cache, single-flight and coalescing routed by key (Part 4).
- Idempotency & Effectively-Once Processing: why a replayed tombstone is safe (Part 5).
- Multi-Region Failover: DNS caches and the tail after a Region moves (Part 5).
- Leases, Fencing Tokens & Distributed Locks: the other meaning of "lease".
- The Retries, Timeouts, Backpressure & Load Shedding loop primitive (coming): why answers arrive after their readers gave up, and how to shed load. The Rate Limiting loop primitive (coming): per-client limits against enumeration.
- Background: Primitive #04: Distributed cache patterns and eviction.
- Drills: The Product Page That Melted Redis (question 1 in Parts 5 and 8, question 2 in Parts 1 and 3) and The Image Service That Made Every Origin Server Sweat (both questions in Part 10).
- Loops: the key-value store (step 2.7), the rate limiter (steps 1.4, 2.4), the URL shortener (steps 1.3, 1.4, 1.5, 2.1, 2.4, 2.6, 3.3, 3.4), the news feed (steps 2.2, 2.5, 3.2), chat (step 2.1), the gaming leaderboard (steps 2.2, 2.4), notifications (step 1.5), search autocomplete (step 2.1), the job scheduler (step 2.1), the digital wallet (step 2.4), hotel reservations (steps 2.1, 3.4), the proximity service (steps 2.1, 2.6), maps and routing (step 1.1), nearby friends (step 1.2), ride-sharing dispatch (step 2.5), YouTube (steps 2.3, 2.4, 3.4), Google Drive (step 2.4), ad-click aggregation (step 2.1), the mobile news feed (step 1.4), mobile stock trading (step 2.6), and the Netflix (step 2.4), Discord (step 3.2), Shopify (step 1.1) and Figma case studies.