You're building a page that shows a user's profile: a name and an avatar. The profile lives in a database, and each time someone opens the page, your app runs a query like SELECT … WHERE id = 42 to fetch it. Say that query takes about 5 milliseconds. For most users that's fine. Then one profile, the one we'll call user:42, belongs to someone famous, and ten thousand people a second want to see it.
The database now answers the same question ten thousand times a second, and it gives the same answer every time, because the profile changes maybe once a week. The obvious move is to keep a copy of the answer somewhere the app can reach faster than the database, and serve that. Writing it takes a dictionary and three lines of code. Living with it is harder, because a copy raises questions the dictionary doesn't answer. How many profiles can you afford to keep? Which one do you throw away when the space runs out? What do readers see after the profile changes? What happens to the database the moment a popular copy disappears? And what if two servers hold different copies?
A copy kept close for speed is called a cache. The app finds things in it by name, and the name it uses for this profile is the string user:42, which we'll call its key. This chapter follows that one key and asks one question the whole way: when we serve the copy instead of asking the database, how sure can we be that it's still right, and what happens to the database when the copy isn't there? We start by simulating a cache to see what a small one buys. Then we work through where it sits, what it forgets, how long a copy may live, and the three ways a cache hurts you.
01How much does a small cache buy?
1.1The counter and the cupboard
A cook keeps salt, oil and a knife on the counter, and the rest of the kitchen in cupboards and a walk-in fridge. The counter is small, so the cook has to decide what earns a spot there. Put the right things on it and most reaching is a step away. Put the wrong things on it and the cook walks to the fridge all day.
A cache is that counter. It's a small, fast store in front of a large, slow one, and the large one, here the database, is called the source of truth because when the two disagree, the database is right. When your app wants a profile, it asks the cache first. If the cache has the key, that's a hit, and the app is done. If not, that's a miss: the app goes to the source of truth and, usually, keeps a copy of the answer in the cache for next time.
The cache has to be smaller than the database, or it would just be a second database. So when it's full and a miss wants to add a new entry, something has to leave. The rule for choosing is a design decision, and the simplest rule is least recently used, or LRU: throw out the entry that nobody has asked for the longest. It's the cook's rule too, since what's been on the counter longest without being touched is the first to go back to the cupboard.
That leaves the question of size. We can't keep every profile, so how much would a cache that holds only a few of them help? It depends on how the requests are spread, and that's easy to find out by simulating.
1.2Trying it: one percent of the keys
The script below generates one million requests over 100,000 profiles. They aren't spread evenly: the most popular profile is twice as likely to be requested as the second, three times as likely as the third, and so on (that's what the weights 1 / (rank + 1) express), so a few profiles get most of the traffic, like our famous user. It then replays the requests through an LRU cache at three sizes: 1%, 5% and 20% of the keys. The two latencies are assumptions, 0.2 ms for an answer from the cache and 5 ms for an answer from the database. For each size it prints the hit rate, the share of requests that were hits, and the average time per request.
Two parts of the code deserve a word. Python's OrderedDict remembers insertion order, and move_to_end marks a key as just used, so the first entry is always the least recently used. popitem(last=False) removes that first entry, which is the eviction. The random generator has a fixed seed, so you get the same table every time.
import bisect, itertools, random
from collections import OrderedDict
KEYS, REQUESTS = 100_000, 1_000_000
HIT_MS, MISS_MS = 0.2, 5.0 # assumed: cache answer vs database answer
rng = random.Random(7)
weights = [1 / (rank + 1) for rank in range(KEYS)] # a few keys are very popular
cum = list(itertools.accumulate(weights))
stream = [bisect.bisect(cum, rng.random() * cum[-1]) for _ in range(REQUESTS)]
def lru_hit_rate(size):
cache, hits = OrderedDict(), 0
for k in stream:
if k in cache:
hits += 1; cache.move_to_end(k)
else:
cache[k] = True
if len(cache) > size: cache.popitem(last=False) # evict least recently used
return hits / len(stream)
print(f"{'cache holds':>12} {'hit rate':>9} {'average request':>16}")
for pct in (1, 5, 20):
h = lru_hit_rate(KEYS * pct // 100)
print(f"{pct:>3}% of keys {h:>9.1%} {h * HIT_MS + (1 - h) * MISS_MS:>13.2f} ms")
print(f"{'no cache':>12} {'0.0%':>9} {MISS_MS:>13.2f} ms") cache holds hit rate average request
1% of keys 50.6% 2.57 ms
5% of keys 66.5% 1.81 ms
20% of keys 80.6% 1.13 ms
no cache 0.0% 5.00 msLook at the first row. A cache holding 1% of the profiles answered half the requests, and that cut the average request from 5.00 ms to 2.57 ms. Twenty times the memory, 20% of the keys, lifted the hit rate to 80.6% and the average to 1.13 ms.
1.3What the table says
Memory has diminishing returns here: the first slice is worth far more than the last. Going from 1% to 5% of the keys raised the hit rate by 16 points, and going from 5% to 20% raised it by only 14 more, for almost four times as much extra memory. That shape comes from the skew in the requests. A few popular keys do most of the work, and a small cache catches them.
Real traffic decides how steep the curve is, and the two latencies here are assumed numbers, so trust the shape of the table more than its exact figures. Notice also how the last column is built: every request costs either 0.2 ms or 5 ms, and the average is a blend of the two. The blend is dominated by the slow part, even when the slow part is rare, and that tells us what to watch when we run a cache for real.
02What a hit rate is worth
2.1Hit rate arithmetic
The last column of the table was a weighted sum, and it's worth writing down because everything else in the chapter leans on it:
average = hit_rate × hit_cost + (1 − hit_rate) × miss_cost
When a miss costs far more than a hit, the second term dominates. So the number to watch is the miss rate, 1 minus the hit rate, because misses are the requests your database sees. Here's what that means for a service doing 10,000 reads a second:
| Hit rate 90%: misses reaching the database | 10,000 reads/s × 0.10 | 1,000 /s |
| Hit rate 95% | 10,000 × 0.05 | 500 /s |
| Hit rate 99% | 10,000 × 0.01 | 100 /s |
| Hit rate drops from 99% to 90% during an incident | 1,000 / 100 | 10× the load |
| Going from 90% to 95% halves the database load | 2× | |
?Why does a five-point drop in hit rate take the database down?
Because the database was sized for the misses it normally sees. At 99% the database sees one read in a hundred. At 90% it sees one in ten, ten times as many, and most databases probably have far less spare capacity (their headroom) than that.
That's the first thing to remember about a cache that protects a database: its hit rate sets the database's load. The next question is where the cache should physically live, because that decides what a hit costs.
2.2What a hit and a miss cost
The 0.2 ms in our simulation was an assumption. To see real figures, the table below times four ways of fetching data from a single Python process, all on one machine:
- A dictionary in the program's own memory. A cache that lives there is called an in-process cache.
- Redis, a separate server that keeps data in memory. The request goes over the loopback interface, network traffic that never leaves the machine, and the time includes decoding the JSON value Redis returns.
- Postgres, asked for one row by its primary key over a Unix socket, a similar machine-local connection. The table fits in
shared_buffers, the area of memory where Postgres keeps recently used table pages, so this lookup is a hit in Postgres's own cache. - Postgres again, running a
GROUP BYthat adds up the whole table. This is the kind of query that does real work on every call.
Each row gives two percentiles. The p50 is the median, the time that half of the reads beat. The p99 is the time that 99 out of 100 reads beat, so it describes the slow end.
| Read path | p50 | p99 |
|---|---|---|
| In-process dict | 0.17 µs | 0.25 µs |
Redis GET + json.loads, loopback | 65 µs | 153 µs |
| Postgres primary-key lookup, Unix socket | 64 µs | 143 µs |
Postgres GROUP BY over the whole table | 141 ms | 278 ms |
GROUP BY), median of three runs.?Why was Redis no faster than Postgres here?
Because the Postgres lookup was already a hit, in Postgres's own buffer cache. Both calls are one round trip, a request sent to another process and an answer sent back, and the round trip is most of the cost. A remote cache only saves latency when the miss path does real work: a join, an aggregate, a call to another service. The GROUP BY row is the kind of thing you'd cache. Notice too that the in-process dictionary is about four hundred times faster than either server, because it has no round trip at all. We'll come back to that in section 10.
Wherever the cache sits, someone has to fill it and keep it in line with the database. Who does that, and when, is the next question.
03Who fills the cache, and what happens on a write
3.1Cache-aside: the app is in charge
The simplest arrangement leaves the cache dumb and puts the app in charge. The app asks the cache for user:42. On a miss it queries the database itself, then stores the answer in the cache with a TTL, a "time to live": a timer, say 60 seconds, after which the cache throws the entry away by itself. (Section 7 is about why that timer matters.) This is called cache-aside, or look-aside, and the Facebook memcache paper describes it as "a demand-filled look-aside cache" (Nishtala et al., NSDI 2013). Step through one page view of user:42, then a change to the profile:
user:42 arrives. The app asks the cache first, by the key it builds from the request. The cache holds another profile but not this one.The last step is the one people find odd. The cache holds a stale copy, so why not overwrite it with the new profile instead of deleting it?
?Why delete the key instead of writing the new value into it?
Take two writers changing the same profile. Writer A saves a new bio and writer B saves another a moment later, so in the database A comes first and B second, and B's bio wins. Each writer then updates the cache as well. Nothing forces those two cache updates to arrive in the same order as the database writes, and if B's update lands first and A's lands second, the cache holds A's older value and nothing will ever correct it until the TTL runs out. Deleting avoids this. Facebook's answer is short: "We choose to delete cached data instead of updating it because deletes are idempotent." An operation is idempotent when doing it twice, or in a different order, leaves the same result, and here two deletes in any order leave the same empty slot. The next reader finds the empty slot, loads whatever the database holds then, and stores it.
3.2Four other arrangements
Cache-aside is one of five common patterns. The others move the filling, and sometimes the writing, from the app into the cache itself. Each one answers two questions: who fills the cache on a miss, and what happens to the cache when the data is written?
| Pattern | On a read miss | On a write | Good for | Watch out for |
|---|---|---|---|---|
| Cache-aside (look-aside) | The application reads the database, then sets the cache | The application writes the database, then deletes the key | Most application caches in front of a database | Races between a slow reader and a writer (section 9) |
| Read-through | The cache library calls a loader, a function you supply that fetches the value from the database | Usually paired with one of the write patterns | In-process caches like LoadingCache in Caffeine, a Java caching library | Loader errors and timeouts now live inside the cache |
| Write-through | As read-through | Writes go to the cache, which writes the database synchronously | Data that's read right after it's written | Every write pays both latencies; the cache fills with data nobody reads |
| Write-behind (write-back) | As read-through | The cache acknowledges, then writes the database later | Absorbing write bursts, counters | The cache holds the only copy for a while. A crash loses acknowledged writes |
| Write-around | As cache-aside | Writes skip the cache entirely; entries age out by TTL | Write-heavy data that's rarely re-read | Readers see old data until the TTL expires |
Write-behind is the only one that changes what happens when the cache is lost. The others treat the cache as disposable, which is what lets you evict from it, restart it and flush it without losing anything. Disposable also means finite, and that raises the question we skipped in section 1: when the cache is full, what should it forget?

04What to forget first
Our simulation used LRU without asking whether it was any good. Every insert into a full cache evicts something, and the rule for choosing the victim is the eviction policy. An eviction policy is a guess about which entry you're least likely to need again.
There is a perfect guess, but nobody can build it. Belady's algorithm, also called OPT, evicts the entry whose next use is furthest in the future. It needs to know the future, so it only works in a simulator, where it sets the ceiling that real policies are compared with. We'll see that ceiling in section 6.
4.1LRU, and why nobody implements it exactly
The textbook LRU is a hash map plus a doubly linked list. The hash map finds an entry by key, and the list orders entries from most to least recently used. A hit moves the entry to the head of the list, and eviction takes the tail. Both operations take constant time, so it looks ideal.

?So what's wrong with a linked list?
Every hit is a write. Moving an entry to the head changes pointers that every other thread also wants to change, so a concurrent LRU needs a lock, a gate that lets one thread through at a time, on every read. On a server with many cores that lock becomes a bottleneck, and it's the reason most production caches approximate LRU instead of implementing it exactly.
4.2CLOCK: an approximation with no list
One of the oldest approximations is CLOCK. The entries sit in a ring, a circular array, and each carries a small counter, or in the simplest version a single reference bit. A hit does one cheap thing: it raises that entry's counter. To evict, a "hand" sweeps around the ring. An entry whose counter is above zero gets its counter lowered and the hand moves on. The first entry the hand finds already at zero is the victim.
Look at the ring. The hand is on A, whose counter is 2, so the sweep lowers it to 1 and moves on. It reaches B, whose counter is already 0, and B is evicted. A was used recently, and the counter let it survive this pass, with no list to reorder.
Postgres's buffer manager, the part that decides which table pages stay in shared_buffers, uses this idea with a counter capped at BM_MAX_USAGE_COUNT = 5. A buffer that's pinned, meaning a query is using it right now, can't be taken at all.
/* Nothing on the freelist, so run the "clock sweep" algorithm */
trycounter = NBuffers;
for (;;)
{
buf = GetBufferDescriptor(ClockSweepTick());
/*
* If the buffer is pinned or has a nonzero usage_count, we cannot use
* it; decrement the usage_count (unless pinned) and keep scanning.
*/
local_buf_state = LockBufHdr(buf);
if (BUF_STATE_GET_REFCOUNT(local_buf_state) == 0)
{
if (BUF_STATE_GET_USAGECOUNT(local_buf_state) != 0)
{
local_buf_state -= BUF_USAGECOUNT_ONE;
trycounter = NBuffers;
}
else
{
/* Found a usable buffer */
/* ... */
return buf;
}
}
/* ... all buffers pinned: error out ... */
UnlockBufHdr(buf, local_buf_state);
}The loop is the sweep. ClockSweepTick() moves the hand, a nonzero usage_count is lowered and skipped, and the first buffer at zero is returned. A hit costs one increment on the small record Postgres keeps for each buffer, with no list to reorder. A buffer touched five times survives five passes of the hand, so the counter adds a little frequency to the recency. Redis approximates LRU differently, by sampling five keys and evicting the oldest of them; chapter 22 covers that.
CLOCK makes hits cheap, but it has the same blind spot as exact LRU, and the next subsection shows what it is.
4.3Where recency fails: scans and loops
LRU assumes that what you touched recently, you'll touch again. Two common access patterns break that.
- A scan. A batch job, a backup or a report reads a large number of keys once each. Under LRU every one becomes "most recent" and pushes out the hot set.
- A loop slightly larger than the cache. The job reads the same keys in the same order, over and over, and the one it needs next is always the one LRU evicted a moment ago.
Here is a scan hitting a tiny cache of four profiles. user:42 and user:7 are the popular ones. A nightly report reads four other profiles, once each, and will never read them again:
user:42 and user:7 are hot, read all day. user:3 and user:9 are cold. The hot keys were touched most recently, so LRU would evict the cold ones first.The loop is worse than the scan, because nothing ever recovers. Try to predict it:
A cache holds 1,000 entries. A job reads keys 0 to 1,099 in order, then starts again, 500,000 reads in total. What hit rate does LRU get?
Databases defend against scans explicitly. When Postgres reads a whole table from start to finish and the table is bigger than a quarter of shared_buffers, it keeps the scan's pages out of the main sweep. The scan reuses a private ring of buffers, 256 KB in all (BAS_BULKREAD in the same freelist.c), so a big report can't evict everyone else's pages. Linux's page cache splits its pages into an active and an inactive list for the same reason: a page normally has to be touched twice before it's promoted to the active list, and pages are reclaimed from the inactive one first.
Those defenses all exist because recency alone can't tell a scan key from a hot key. The other obvious signal is how often a key is used.
4.4LFU, and why frequency goes stale
Least frequently used, or LFU, evicts the entry with the fewest hits. It shrugs off scans, because a key read once has a count of one and leaves first.
?Then why isn't LFU the default everywhere?
Two reasons. The first is memory: exact LFU needs a counter for every key, including keys that have left the cache, or a returning key starts from zero. The second is staleness. Counts only go up, so yesterday's popular keys keep today's newcomers out.
You can see the second one in a simulation. When the trace's hot set changes halfway through, an LFU that counts only while a key is cached falls from 44% to 25.5%, well below LRU's 34%. Old keys have counts in the thousands, and the new ones can't catch up before they're evicted.
Any usable frequency policy therefore needs aging: counts that decay, so popularity has a half-life. Redis's LFU uses a logarithmic 8-bit counter that decays by lfu-decay-time. W-TinyLFU, the next section, halves every counter periodically.
So recency is fooled by scans and plain frequency is fooled by change. To get both, we need a different kind of question than "who leaves?".
05Admission: deciding who gets in
Every policy so far answers "which entry leaves?" The scan scene shows another way to look at it. The scan keys were read once and will never be read again, and they still got in and threw out the hot keys. What if the cache could refuse them at the door? Admission is the decision, made about each newcomer when the cache is full, whether it's worth more than the entry it would replace. If it isn't, the cache drops the newcomer and changes nothing.
5.1The gate
Start with the simple version. Count how often each key is requested. When a newcomer arrives and the cache is full, compare its count with the count of the entry that would be evicted, the victim, and keep whichever is more popular. A scan's keys all have a count of one, so they lose against any key that's been requested twice, and the hot set survives. This is the idea behind TinyLFU.
It has two problems. To compare a newcomer with the victim, you need the newcomer's count, and the newcomer isn't in the cache yet, so you have to keep counts for keys the cache doesn't hold, which means every key ever seen. That costs more memory than the cache itself. The second problem is the one from LFU: counts that only go up go stale. A third problem turns up once we've solved those two, and we'll meet it in section 5.3.
5.2The sketch: 4-bit counters that halve themselves
The fix for the memory problem is to count approximately. Caffeine, a widely used Java cache library built on this idea, keeps a fixed-size table of small counters, and every key is hashed four ways to pick four of them. This structure is called a count-min sketch. An increment bumps all four, and the estimate of a key's count is the smallest of its four. Other keys share those counters, so collisions can only inflate an estimate, and taking the minimum limits the damage. The sketch sees every access, hit or miss, so it remembers keys the cache doesn't hold, in a few bits each.
public void increment(E e) {
/* ... */
int[] index = new int[8];
int blockHash = spread(e.hashCode());
int counterHash = rehash(blockHash);
int block = (blockHash & blockMask) << 3;
for (int i = 0; i < 4; i++) {
int h = counterHash >>> (i << 3);
index[i] = (h >>> 1) & 15;
int offset = h & 1;
index[i + 4] = block + offset + (i << 1);
}
boolean added =
incrementAt(index[4], index[0])
| incrementAt(index[5], index[1])
| incrementAt(index[6], index[2])
| incrementAt(index[7], index[3]);
if (added && (++size == sampleSize)) {
reset();
}
}
/** Reduces every counter by half of its original value. */
void reset() {
int count = 0;
for (int i = 0; i < table.length; i++) {
count += Long.bitCount(table[i] & ONE_MASK);
table[i] = (table[i] >>> 1) & RESET_MASK;
}
size = (size - (count >>> 2)) >>> 1;
}increment computes four counter positions for the key and bumps each one. Once enough increments have happened, reset() halves every counter in the table. Four details of the code are there on purpose:
| Detail | What it does |
|---|---|
Sixteen 4-bit counters per long | A counter saturates at 15. You only need to know "more popular than the victim", not the exact count |
| All four counters in one 64-byte block | The CPU fetches memory in 64-byte pieces, so a lookup touches one piece of memory instead of four, as the class comment explains |
sampleSize = 10 × maximum | After ten increments per cache slot, reset() halves every counter. That's the aging that section 4.4 said LFU needs |
table sized to the cache's maximum, rounded to a power of two | One long per entry: roughly 8 MB of sketch for a million-entry cache |
Halving is what fixes LFU's staleness. A key that was hot an hour ago loses half its count every sampling period, so new popularity wins within a few periods. In the "hot set replaced" trace that we'll see in section 6, the full design built on this sketch lost almost nothing.
5.3A window for newcomers, and a main region
The third problem is that pure frequency admission is too harsh on a burst. Think of a profile that's about to be read a hundred times in the next second. On its first read its count is one, so a strict gate would reject it before it had a chance to show what it was worth.
?Why is the window there, if the sketch already filters newcomers?
Because every newcomer needs somewhere to live while it proves itself. A burst key that is rejected on read one never gets to read two. So the design puts a small window in front of the gate: a tiny LRU, 1% of the cache at the start, that every new entry enters. A key that's part of a burst collects hits and a count there, without disturbing the rest of the cache. TinyLFU with this window added is called W-TinyLFU, and it's the policy Caffeine implements.
The other 99% of the cache is the main region, where the gate applies. Inside it are two LRU lists. New arrivals join the first, the probation list, which takes 20% of the main region. A hit on an entry in probation promotes it to the second, the protected list, which takes the other 80%. When protected overflows, its oldest entry drops back to probation, so it gets one more chance before leaving the cache. Two LRU lists arranged like this are called a segmented LRU, or SLRU.
When the window overflows, its oldest entry becomes the candidate. The victim is the oldest entry in probation. The sketch estimates both counts, and only the more frequent one stays. Here's one such duel, drawn with one slot in the window so it fits on the page. The number under each key is what the sketch currently says about it:
user:42 has the highest count and sits in protected. user:9 is the oldest entry in probation, so it's the next victim. user:98 is in the window: a key the nightly report read once.A candidate that wins looks the same in reverse. A key that has been read four times while it sat in the window arrives at the gate with a count of 4, beats the victim's 2, and takes the victim's place in probation. The victim is evicted.
Putting the parts together, here is what happens to one new key arriving at a full W-TinyLFU cache:
| Step | What happens |
|---|---|
| Access | A read misses on key K. The loader fetches it, and the cache stores it |
| Sketch | The access increments K's counters. The sketch sees every access, hit or miss |
| Window | K goes into the window. If K is part of a burst, it can get hits here |
| Gate | When the window overflows, its oldest entry, the candidate, meets the victim from probation. The sketch estimates both counts |
| Outcome | The higher estimate stays. If the candidate loses, it's dropped and the main region doesn't change |
| Promotion | An admitted entry lives in probation, first in line to be evicted unless it's used again. A hit in probation promotes it to protected |
5.4Tuning itself, and defending against attacks
The 1% window suits frequency-heavy workloads, but it's too small for recency-heavy ones. Caffeine adjusts it at run time with a hill climber. It moves the split between window and main by a step (starting at 6.25% of capacity, decaying by 0.98 each time), keeps going while the hit rate improves, reverses when it drops, and restarts if the hit rate shifts by more than 5%. Those are the HILL_CLIMBER_* constants in BoundedLocalCache.java.
Admission also opens an attack. If an attacker can make the victim look popular, for instance by crafting keys that collide in the sketch, no new entry is ever admitted. Caffeine's admit() answers with a little randomness:
boolean admit(K candidateKey, K victimKey) {
int victimFreq = frequencySketch().frequency(victimKey);
int candidateFreq = frequencySketch().frequency(candidateKey);
if (candidateFreq > victimFreq) {
return true;
} else if (candidateFreq >= ADMIT_HASHDOS_THRESHOLD) {
// The maximum frequency is 15 and halved to 7 after a reset to age the history. An attack
// exploits that a hot candidate is rejected in favor of a hot victim. The threshold of a warm
// candidate reduces the number of random acceptances to minimize the impact on the hit rate.
int random = ThreadLocalRandom.current().nextInt();
return ((random & 127) == 0);
}
return false;
}A warm candidate (count of 6 or more) that loses still gets in one time in 128. That's rare enough not to hurt the hit rate and often enough that a pinned victim can't hold its slot for long.
W-TinyLFU is one answer among several. Which of them wins depends on the traffic, so the next step is to line them up and run them on the same requests.
06How the policies compare
6.1The field
Most modern policies combine recency and frequency somehow. Now that we've seen how each ingredient works, here is the field in one place:
| Policy | Idea | Used in |
|---|---|---|
| LRU | Evict the least recently used | Most libraries' default; the baseline |
| CLOCK | LRU approximation with a reference bit or counter and a sweeping hand | Postgres buffers, OS page replacement |
| 2Q / SLRU | New entries go to a probation area; a second hit promotes them | Postgres 8.0.2 briefly; the main region of W-TinyLFU |
| ARC (Megiddo and Modha, FAST 2003) | Balances a recency list and a frequency list, adapting the split using "ghost" lists: lists that remember the keys of recent evictions, without their data | ZFS. Postgres shipped it in 8.0 and replaced it in 8.0.2 to avoid an IBM patent (release notes) |
| W-TinyLFU (Einziger, Friedman and Manes) | An LRU window, an SLRU main region, and a frequency sketch that decides admission | Caffeine (Java), Ristretto (Go) |
| S3-FIFO (Yang et al., SOSP 2023) | A small FIFO (10%) filters out keys read only once before a main FIFO (90%); a ghost queue remembers recent evictions | Newer; lock-free friendly because FIFOs need no reordering on a hit |
S3-FIFO takes a different route to the same goal as the admission gate. A FIFO is a first-in, first-out queue, and unlike an LRU list it never reorders entries on a hit. New keys enter the small queue, and only the ones that are requested again while they're there move on to the main queue. A key requested only once, a one-hit wonder, falls off the end of the small queue without reaching the main one. That matters because one-hit wonders are common: across 6,594 production traces, the S3-FIFO authors found that the median share of objects accessed only once was 26%, and much higher over short windows (blog summary). Every one-hit wonder kept out of the main queue leaves its slot to a key that will be asked for again.
6.2A comparison on synthetic traces
To compare the policies on equal terms, the table below runs five of them on the same simulated requests. LRU and in-cache LFU are textbook. S3-FIFO follows the paper's description, and W-TinyLFU follows Caffeine's constants without the hill climber. OPT is the Belady ceiling from the start of section 4. Four synthetic traces were used. Three of them have 2 million requests drawn from 100,000 keys with a Zipf skew of 0.9, which is a pattern where the key of popularity rank n is requested in proportion to 1/n^0.9, so a few keys dominate, as in our first simulation. Those three are a stable one, one with a 2,000-key scan every 10,000 reads, and one where the hot set is replaced halfway through. The fourth is the loop over 1,100 keys from the prediction above. The cache holds 1,000 entries (1% of the keys) in every row except the last, which repeats the stable trace with a cache of 5,000.
| Trace | LRU | LFU | S3-FIFO | W-TinyLFU | OPT (ceiling) |
|---|---|---|---|---|---|
| Zipf 0.9, stable | 34.2% | 44.3% | 45.1% | 45.1% | 53.3% |
| Zipf + a 2,000-key scan every 10,000 reads | 33.2% | 44.2% | 45.1% | 44.9% | 53.3% |
| Hot set replaced halfway | 34.2% | 25.5% | 44.9% | 44.9% | 53.3% |
| Loop over 1,100 keys | 0.0% | 0.0% | 76.1% | 79.3% | 90.7% |
| Zipf 0.9, stable, cache of 5,000 | 51.5% | 59.3% | 60.0% | 60.5% | 69.6% |
Read it this way. On a skewed workload, the gap between LRU and a frequency-aware policy is roughly ten points of hit rate. At a cache of 1,000 that's a miss rate of 66% against 55%, a sixth less backend load. Plain LFU matches the modern policies until popularity moves, then falls apart, which is the staleness of section 4.4. S3-FIFO and W-TinyLFU stay close to each other and to the ceiling on every trace.
The scan row is smaller than you might expect: on this trace, scans cost LRU roughly a point. With Zipf skew, the hot keys come back within a few thousand requests, so the damage from each scan is short-lived. Scans hurt LRU most when misses are expensive and the hot set is large, which is the buffer-pool case Postgres defends against.
So far we've treated the cache as a place that is only ever too small. It also holds copies that grow old, which is a separate problem, and a timer is the first defence.
07How long may a copy live?
Eviction is about space. Expiry is about time: a TTL says how long a copy may be served without checking the source. Back in the cache-aside scene we set user:42 with a TTL of 60 seconds. We'll see now what that number promises.
7.1What a TTL promises
A TTL is an upper bound on staleness, if nothing else invalidates the entry. With a 60-second TTL and no invalidation, a user can see a profile up to 60 seconds old. With invalidation, which is the delete on write from section 3, the TTL becomes the backstop for the invalidations you'll eventually lose to a crash, a bug or a dropped message.
?Why keep a TTL at all, if every write deletes the key?
Because deletes get lost. A process can crash between committing to the database and deleting the key. A network partition can swallow the delete. A new code path can write the table and forget the cache. Without a TTL, each of those leaves a wrong value in the cache forever. With one, it's wrong for a bounded time.
7.2Three habits for TTLs
The TTL on a single key is easy. The trouble starts when many keys share a TTL, or when a TTL is the only thing that stops a bad request. Three habits help. The first adds jitter, a small random amount added to or taken off each TTL, so that entries written at the same moment expire at slightly different moments.
| Habit | Why |
|---|---|
| Add jitter: TTL × random(0.9, 1.1) | Keys written together (a warm-up job, a deploy, a bulk import) otherwise expire together, and the misses arrive as one burst |
| Cache misses too (negative caching), with a short TTL | A lookup for a key that doesn't exist otherwise hits the database every time. That's also the easiest way to attack a cache: request random IDs |
| Separate soft and hard expiry | Past the soft TTL, serve the old value and refresh in the background. Past the hard TTL, block and reload. This is stale-while-revalidate (section 8.4) |
The burst of misses that jitter spreads out has a name, and it's the first of the three ways a cache hurts you. Jitter can't help when the burst comes from a single key, and a single hot key like ours is where it's worst.
08Stampedes: the instant a hot key expires
8.1What a stampede looks like
Go back to user:42, read thousands of times a second. When its copy expires, every request that arrives in the next moments finds the slot empty. Each one goes to the database and runs the same expensive query, and each one stores the same value when it's done:
user:42. The copy is in the cache with a second left on its TTL, so every read is a hit and the database sees nothing.A small test shows how big the pile gets. Eight processes with eight threads each read a single hot key from Redis in a loop, with a 5 ms pause between one read and the next. The value takes 50 ms to compute and lives for 2 seconds. Over a 10-second run the key expires four times. With naive cache-aside the run made 256 backend calls (queries to the database), which is 64 per expiry, one for every reader thread, all at the same moment. Scale that to a few hundred servers and a query that takes a second, and the database sees a burst of hundreds of identical queries every time the TTL fires.
?Why did every thread miss, not just the first?
Because the key stays missing for the whole 50 ms the first reader spends computing it. Every read that lands in that window also misses, and at this rate that's every thread. A longer recompute means a bigger herd.
How can we stop 64 threads from asking the same question? The first idea is to make them share one answer.
8.2Coalescing: one fetch per key
A first fix is to make concurrent misses for the same key share one fetch. Go's singleflight is the smallest version of the idea:
// Do executes and returns the results of the given function, making
// sure that only one execution is in-flight for a given key at a
// time. If a duplicate comes in, the duplicate caller waits for the
// original to complete and receives the same results.
// The return value shared indicates whether v was given to multiple callers.
func (g *Group) Do(key string, fn func() (interface{}, error)) (v interface{}, err error, shared bool) {
g.mu.Lock()
if g.m == nil {
g.m = make(map[string]*call)
}
if c, ok := g.m[key]; ok {
c.dups++
g.mu.Unlock()
c.wg.Wait()
/* ... re-panic or Goexit if the leader did ... */
return c.val, c.err, true
}
c := new(call)
c.wg.Add(1)
g.m[key] = c
g.mu.Unlock()
g.doCall(c, key, fn)
return c.val, c.err, c.dups > 0
}Whoever calls first for a key registers a call and runs fn. Everyone who arrives while it's running waits on the same WaitGroup and gets the same result. Caffeine's LoadingCache does the same per key inside the cache.
In the same test, singleflight cut the 256 backend calls to 32, which is 8 per expiry.
?Why 8 calls per expiry, and not one?
Because singleflight only coalesces within one process. The test had eight processes, so each expiry had eight leaders, one per process. That's a big improvement, 64 down to 8, but the database still sees one query per server on every expiry. To get to one per fleet, the coordination has to live somewhere shared: the cache itself.
A lock in Redis (SET lock:K 1 NX PX 1000, which sets the lock key only if it doesn't already exist, and lets it expire after a second) does that, and cut the calls to 6 for 4 expiries. The two extra calls came from a loser that saw the miss, then won the lock a moment after the winner had filled the key and released it. Check the cache again after taking the lock and those go away. Losers also have to wait for the winner, and checking back on a timer gave them the slowest reads in the comparison table of section 8.5. And the lock needs a timeout, in case the winner dies holding it.
8.3Leases: Facebook's version
Facebook's memcache built the lock into the cache server and called it a lease. On a miss, memcached hands the client a 64-bit token bound to that key. Only a SET carrying a valid token is accepted.
On a set of keys prone to thundering herds, Facebook measured the peak database query rate from cache misses at 17K/s without leases and 1.3K/s with them, over one week (NSDI 2013, section 3.2.1). That same token also fixes the stale-set race in section 9, because a delete invalidates any outstanding token.
8.4Refreshing before expiry: XFetch and stale-while-revalidate
Both fixes so far still make readers wait during the recompute. The other family avoids the miss entirely, by refreshing before the key expires.
Probabilistic early expiration has each reader decide, at random, to recompute a little early, with the odds rising as expiry approaches. Vattani, Chierichetti and Lowenstein proved the optimal shape and called it XFetch (VLDB 2015):
function XFetch(key, ttl; β = 1)
value, Δ, expiry ← CacheRead(key)
if !value or Time() − Δ·β·log(rand()) ≥ expiry then
start ← Time()
value ← RecomputeValue()
Δ ← Time() − start
CacheWrite(key, (value, Δ), ttl)
return valueΔ is roughly how long the last recompute took, so slow values refresh earlier. rand() returns a number between 0 and 1, so log(rand()) is negative, and subtracting it pushes each reader's idea of "now" forward by a random, exponentially distributed amount. Most readers see a fresh value. A few see it as expired shortly before it is, and one of them usually refreshes it first. There's no lock and no coordination, only two extra numbers, Δ and the expiry time, stored with the value.
In the stampede test, XFetch made 10 or 11 backend calls over the four expiries, two or three each.
?Why more than one call per expiry?
Because with sixty-four readers polling every few milliseconds, several of them roll "expired" within the same 50 ms recompute. The paper's β = 1 default trades some duplicate work for zero coordination. Raising β refreshes earlier and cuts the overlap. Even so, it went from 64 calls per expiry to under 3, with readers never blocked.
Stale-while-revalidate is the version HTTP caches use (RFC 5861). Past a soft deadline, the cache keeps serving the old value while exactly one refresher (holding a lock) fetches the new one. Its slowest reads were the fastest of the five in the test, because no reader ever waited for the backend once the key existed.
8.5The five side by side
Now that each strategy has been described, here are all five from the same test, with the counts over 10 seconds and four expiries. The last column, p999, is the read time that 999 out of 1,000 reads beat, so it shows the very slowest readers.
| Strategy | Backend calls in 10 s | Most at once | Reader p999 (median of 3) |
|---|---|---|---|
| Naive cache-aside | 256 | 64 | 52 ms |
| singleflight, per process | 32 | 8 | 52 ms |
| Redis lock (SET NX), losers poll every 10 ms | 6 | 1 | 58 ms |
| XFetch, β = 1 | 10–11 | 3–4 | 36 ms |
| Stale-while-revalidate, one refresher | 5–6 | 1 | 24 ms |
And as a decision table:
| Technique | Coordination | Readers wait? | Cost |
|---|---|---|---|
Coalescing (singleflight, LoadingCache) | In-process | Yes, for the one fetch | One fetch per process per expiry |
| Lock or lease in the cache | Shared | Yes, and they poll | A lock timeout to tune; the holder can die |
| XFetch | None | Rarely | A few duplicate fetches; store Δ and expiry with the value |
| Stale-while-revalidate | A lock for the refresher | No | Serves data past its TTL; needs a hard TTL too |
A stampede is a copy that has disappeared. The opposite problem is a copy that's still there and wrong, and it can happen even when every write remembers to delete the key.
09Invalidation: keeping the copy honest
Once data changes, every cached copy of it is wrong until someone removes it. Sending the delete is easy. The hard part is the ordering between the delete and a reader that was already halfway through filling the cache.
9.1The stale-set race
This is the classic cache-aside race, and it needs only one slow reader and one writer. Go back to user:42, written K below, with its avatar about to change from a.png (call that value v1) to b.png (v2):
?Why doesn't deleting twice fix it?
A common patch is the "delayed double delete": the writer deletes the key, waits a moment, and deletes it again. That catches a stalled reader only if the reader wakes up and sets its old value before the second delete. Nobody can promise how long a garbage-collection pause or a network stall will last, so the patch makes the race rarer and leaves it possible.
9.2Closing the race
A race like this closes when the cache can tell that a set is based on data older than a delete it has already seen. There are three ways to give it that information.
| Approach | How it works | Where you see it |
|---|---|---|
| Leases | A miss issues a token; a delete invalidates it; a set with an invalid token is refused | memcache at Facebook (section 8.3), which the paper compares to load-link/store-conditional, a CPU instruction pair where the store fails if the memory changed in between |
| Versioned values | Store the row's version (an updated_at or a counter) with the value, and only set if it's newer than what's there, using a Lua script (which Redis runs without interruption) or WATCH (which cancels the set if the key changed meanwhile) | Any cache where rows carry a version column |
| Invalidate from the commit log | A daemon follows the database's replication log, its ordered record of every committed change, and issues deletes after each commit, in commit order | Facebook's mcsqueal; Debezium or another change-data-capture pipeline (a tool that turns a database's change log into a stream of events) elsewhere |
Facebook's mcsqueal runs on every database, reads the SQL statements that commit, extracts the keys to delete, and broadcasts the deletes to every front-end cluster in the region. Because it reads the log, a delete can't be lost to an application crash between commit and delete. Facebook also notes that "only 4% of all deletes issued result in the actual invalidation of cached data": most of the time the key wasn't cached anyway.
Everything so far has assumed one cache. Real services run on many servers, and the fastest cache of all lives inside each of them.
10Coherence across a fleet
The fastest cache is in your process: section 2.2 timed it at 0.17 µs against 65 µs for Redis. But with 200 application servers, there are 200 copies of every hot key, and a write has to reach all of them. Keeping many copies in agreement is called coherence, the same word hardware uses for CPU caches (chapter 02).
10.1Four ways to keep local copies in line
| Approach | Staleness | Cost | Fails how |
|---|---|---|---|
| Short TTL only (seconds) | Up to one TTL | Nothing to build; more misses | Wrong for up to a TTL after every write |
| Broadcast invalidations (Redis pub/sub, where servers publish messages to a channel that others subscribe to; or Kafka) | One message delay | A publish per write, a subscription per server | A server that misses a message stays stale until the TTL |
Server-assisted tracking (Redis CLIENT TRACKING, Redis 6+) | One message delay | Server memory per tracked key and client | As above, plus forced invalidations when the tracking table fills |
Versioned keys (user:42:v17) | None for readers who know the version | An extra lookup to learn the current version | The version lookup is itself a cache problem |
Redis's tracking has two modes (docs). In the default mode, the server remembers which clients read which keys and sends each client invalidations only for its keys. In broadcasting mode (BCAST PREFIX user:), it remembers nothing and sends every change under the prefix to every subscriber.
?Why would Redis invalidate a key that didn't change?
Because the table of who read what, the invalidation table, has a size limit (tracking-table-max-keys). When it's full, Redis evicts an entry "by pretending that such key was modified", and sends the invalidation anyway. The client loses a good cached copy, but it never keeps a stale one. That's the right way for a coherence mechanism to fail when it runs out of memory: the price is extra misses, which are slow but correct.
10.2The races that come with local caches
Local caches bring the race from section 9 back in a new form. The Redis docs describe it for a client that reads data on one connection, marked [D] below, and receives invalidations on another, marked [I]:
[D] client -> server: GET foo
[I] server -> client: Invalidate foo (somebody else touched it)
[D] server -> client: "bar" (the reply of "GET foo")Here the invalidation overtakes the reply, so the client caches a value that has already been declared stale. Their fix is a placeholder: mark foo as "caching in progress" before sending the GET, let an invalidation delete the placeholder, and don't store the reply if the placeholder is gone. It's more or less a local version of Facebook's lease.
A second rule from the same docs: if you lose the invalidation connection, flush the local cache. You can't know what you missed. Ping the invalidation channel periodically, and treat a silent channel as a lost one.
10.3Hot keys and failed cache servers
A distributed cache spreads keys across servers by hashing, usually consistent hashing, a scheme in which one server's failure moves only its share of the keys. Two things break that tidy picture, and user:42 is the cause of the first.

Hot keys. Hashing spreads keys, not load. Facebook noted that "a single key can account for 20% of a server's requests." A celebrity's profile or a homepage config lands on one cache server, and that server saturates while the others idle. Usual fixes are a small local cache in front for the hottest keys, or storing N copies under K#1 … K#N and reading a random one.
Failed servers. When a cache server dies, the obvious fix is to rehash its keys onto the survivors. Facebook deliberately didn't do that, because a hot key could then overload the server it moved to, and so on around the ring of servers that the hashing lays out. Instead, failed requests go to Gutter, a pool of idle servers about 1% of the cluster's size, whose entries expire quickly. That turned 10–25% of failures into hits and cut client-visible failures by 99% (NSDI 2013, section 3.3).
With the three failures covered, what remains is knowing which of them is happening to you.
11Operating a cache
11.1What to measure
Each failure we met shows up in a different number, so these are the ones to graph.
- Hit rate per key family. A global number hides a 20% hit rate on the one prefix that matters (section 2).
- Misses per second reaching the source. That's the load the database is provisioned for.
- Evictions against expirations. Lots of evictions means the cache is smaller than the working set (section 4); lots of expirations at once points at a stampede (section 8).
- Recompute time of expensive values, since it sets the stampede window.
On Redis, the first three are fields in one command, and the same names exist in memcached's stats (get_hits, get_misses, evictions):
# Hit rate and misses per second reaching the source (section 2): sample twice and subtract
redis-cli INFO stats | grep -E 'keyspace_hits|keyspace_misses'
# Evictions against expirations (sections 4 and 7)
redis-cli INFO stats | grep -E 'evicted_keys|expired_keys'
# Which keys are hot? (section 10.3) needs an LFU eviction policy
redis-cli CONFIG SET maxmemory-policy allkeys-lfu # changes eviction for the whole server
redis-cli --hotkeys
redis-cli OBJECT FREQ user:42 # the logarithmic access counter for one key11.2Rules that hold up
- Delete on write, and keep a TTL as well. The delete handles the normal case and the TTL bounds every case you forgot (sections 3 and 7).
- Size the database for the misses. A cold or failed cache sends the whole read load to it (section 2).
- Jitter every TTL on keys that are written together.
- Guard every hot key against a stampede with coalescing, a lease, or early refresh (section 8).
- Use a scan-resistant policy if batch jobs share the cache with real traffic (sections 4 and 5).
- Give every local cache an invalidation path and a flush on disconnect (section 10).
11.3What you trade for what
| You get | You pay | When the bill arrives |
|---|---|---|
| Reads served in microseconds instead of milliseconds | Data that can be stale for up to a TTL | When a user sees an old profile |
| A database sized for the misses, not the reads | The database's capacity now depends on the hit rate | On a cold start or a hit-rate drop |
| A scan-resistant policy (W-TinyLFU) | More machinery than LRU: a sketch of about 8 MB for a million entries, plus an aging step | When you write it yourself instead of using a library |
| A lock or lease against stampedes | A timeout to tune, and a holder that can die | When the winner crashes mid-fetch |
| Fast in-process copies | Many copies that must be kept coherent | As different servers showing different values |
11.4Symptom, cause, fix
| Symptom | Likely cause | Fix |
|---|---|---|
| Database load spikes on a regular beat | Hot keys expiring together, or one key stampeding | Jittered TTLs, coalescing, XFetch or stale-while-revalidate |
| Database overloaded after a cache restart or deploy | Cold cache; nothing to coalesce | Warm from a peer, ramp traffic, rate-limit misses at the source |
| Hit rate falls after a nightly job | A scan evicting the hot set under LRU | A scan-resistant policy (W-TinyLFU, S3-FIFO, Redis allkeys-lfu), or bypass the cache for batch reads |
| Users occasionally see old data that never goes away | The stale-set race, and no TTL | Add a TTL now; then leases, versioned sets or log-driven invalidation |
| Different servers show different values | Local caches with no invalidation | Short TTL, broadcast invalidation, or Redis client tracking |
| One cache node at 100% CPU, others idle | A hot key | Local cache for that key, or replicate it under several names |
| Lots of lookups for IDs that don't exist | No negative caching; possibly an enumeration attack | Cache "not found" with a short TTL |
| Hit rate fine, latency not improved | Caching something that was already cheap | Cache the expensive computation, or move the cache in process |
12Summary
- A small cache catches a lot, because requests are skewed. In the simulation, a cache of 1% of the keys answered half the requests, and 20% of the keys answered 80.6%.
- The miss rate is the number to watch, because misses are the database's load. Going from 90% to 95% halves the load on the database; dropping from 99% to 90% multiplies it by ten.
- A remote cache is a network hop. In the test, Redis on loopback cost the same as a Postgres primary-key lookup; cache remotely only what's expensive to miss.
- Delete on write, don't update. Deletes are idempotent; out-of-order sets leave old values behind.
- LRU is simple and fragile. It's defeated by loops and hurt by scans, and a concurrent LRU writes on every read.
- Frequency needs aging. Plain LFU collapsed from 44% to 25.5% when the hot set changed; periodic halving fixes it.
- W-TinyLFU decides admission, not just eviction. A newcomer only gets in if the sketch says it's more popular than the victim, and a small window gives bursts time to prove themselves.
- Every TTL is a staleness budget and a safety net. Keep one even when you invalidate, because invalidations get lost.
- A stampede is one miss multiplied by every concurrent reader. Naive cache-aside made 64 backend calls per expiry in the test. Coalesce in-process, coordinate in the cache, or refresh early.
- A reader that stalls can leave a stale value behind. Leases, versioned sets or log-driven invalidation close the race; a delayed second delete only narrows it.
- Local caches need an invalidation path and a flush on disconnect. Otherwise every server keeps its own version of the truth.
13Build this
A cache policy simulator, and a stampede you can watch.
- Implement LRU, LFU with periodic halving, S3-FIFO and W-TinyLFU against one interface:
access(key) -> hit. Add Belady's OPT from a precomputed next-use array, as the ceiling. - Generate a Zipf trace, then add scans, a loop and a popularity shift, and reproduce the table in section 6. Then find a public trace (the S3-FIFO repository links several) and see whether the ranking holds.
- Put Redis in front of a function that sleeps 50 ms. Run 64 reader threads against one key with a 2-second TTL, count backend calls per expiry, and add singleflight, a
SET NXlock, XFetch and stale-while-revalidate one at a time. - Finally, reproduce the stale-set race of section 9.1 deterministically by adding a sleep between the reader's
SELECTand itsSET. Then close it with a version check.
14Interview questions
beginnerWhat's the difference between cache-aside and write-through?›
In cache-aside, the application owns the cache: on a miss it reads the database and sets the key, and on a write it updates the database and deletes the key. In write-through, writes go through the cache, which writes the database synchronously, so the cache always holds the latest written value. It costs both latencies on every write and fills the cache with data that may never be read.
beginnerWhy do you delete a cache key on write instead of setting the new value?›
Because deletes are idempotent and sets aren't. Two concurrent writers can reach the cache in the opposite order to the database, and a set-on-write leaves the older value cached. Two deletes in any order leave the key empty, and the next reader loads the current row. Facebook's memcache paper gives exactly this reason.
intermediateWhat's a cache stampede, and how do you prevent one?›
A popular key expires, and every concurrent reader misses and recomputes it at once. In a test with 64 readers and a 50 ms recompute, naive cache-aside made 64 backend calls per expiry. Fixes: coalesce misses in-process (singleflight, a loading cache); coordinate through the cache with a lock or lease so one client refills; refresh early at random (XFetch); or serve stale while one refresher runs (stale-while-revalidate). Also jitter TTLs, so keys written together don't expire together.
intermediateWhy does LRU perform badly on a scan, and what do real systems do about it?›
Every key the scan touches becomes most-recently-used, and each one pushes an entry of the hot set toward eviction. A loop slightly bigger than the cache is worse: LRU scored 0% on one in simulation. Postgres runs large sequential scans through a private 256 KB buffer ring. Linux's page reclaim promotes a page to its active list only on a second touch. Caffeine's W-TinyLFU rejects newcomers less popular than the victim, and S3-FIFO filters one-hit wonders through a small FIFO.
intermediateHow does W-TinyLFU work?›
New entries enter a small LRU window, about 1% of capacity. When the window overflows, its oldest entry becomes a candidate, and the main region's eviction victim is the oldest in probation. A count-min sketch with 4-bit counters, which sees every access including misses, estimates both frequencies, and only the more frequent one stays. Counters are halved after ten increments per cache slot, so popularity ages. Caffeine tunes the window size with a hill climber and adds a small random admission against hash collision attacks.
deepDescribe the stale-set race in cache-aside, and two ways to close it.›
A reader misses, reads v1 from the database, and stalls. A writer commits v2 and deletes the key, which isn't cached yet. The reader wakes and sets v1, and it stays there until the TTL. A delayed double delete narrows the window but doesn't close it. To close it: leases, where a miss issues a token, a delete invalidates it, and a set with a dead token is refused; or versioned sets, where the value carries the row version and a set only succeeds if it's newer than what's cached.
deepYou add an in-process cache to 200 servers. How do you keep them coherent, and what can go wrong?›
Options: a short TTL alone; broadcasting invalidations over pub/sub or Kafka; Redis client tracking, per-key or broadcast by prefix; or versioned keys. What can go wrong: an invalidation can overtake the reply it's meant to invalidate if they travel on different connections, which a placeholder entry fixes. A lost subscription leaves a server stale, so flush on disconnect and ping the channel. And a full tracking table forces invalidations, costing hit rate, not correctness. Keep a TTL as the backstop in every case.
deepA cache server dies. Why might rehashing its keys onto the others make things worse?›
Load isn't spread evenly across keys. A hot key can carry a large share of one server's traffic (Facebook cites 20%), and moving it onto a neighbour can overload that neighbour, which fails in turn. Facebook sends failed requests to a small Gutter pool, about 1% of the cluster, with short-lived entries. That converted 10–25% of failures into hits and cut client-visible failures by 99%.
15Go deeper
A 10,000 reads/s service's hit rate drops from 99% to 95%. How much more load does the database see?›
Five times as much: misses go from 100/s to 500/s.
In Caffeine's frequency sketch, what's the largest count a key can have?›
- Counters are 4 bits, and they're halved every 10 × maximum-size increments, so an old count of 15 becomes 7.
Why did per-process singleflight make 8 backend calls per expiry in the test?›
There were eight processes. singleflight coalesces within one process, so each process elected its own leader.
In XFetch, what is Δ and why is it stored with the value?›
How long the last recompute took. Slow values need to refresh earlier, and the random early-expiry gap is scaled by Δ.
Beyond Physical Memory: Policies. The same eviction question for an operating system's page cache: OPT, FIFO, LRU and CLOCK, with a looping workload that LRU handles badly. Free online at ostep.org.
Look-aside caching at scale: leases, Gutter, mcsqueal, cold-cluster warm-up and cross-region consistency. PDF.
The W-TinyLFU design and its evaluation on real traces. arXiv 1512.00727, later in ACM Transactions on Storage.
The proof behind XFetch, and why the exponential gap is optimal. PDF.
S3-FIFO, one-hit wonders, and 6,594 production traces. PDF.
The best-documented production cache policy you can read. Start at
admit() and FrequencySketch, then the hill climber.
v3.1.8.
Tracking, broadcasting, the invalidation table, and the two races every near-cache must handle. redis.io.
The adaptive policy with ghost lists that most later work compares against. PDF.
stale-while-revalidate and stale-if-error, the HTTP names for serving
old data on purpose. rfc-editor.org.
16Related chapters
Sampled LRU and the logarithmic LFU counter, expiry, and why a volatile-*
policy can refuse writes. Chapter 22.
The same problems in hardware: cache lines, eviction and MESI keeping many copies coherent. Chapter 02.
Why a stampede's burst of identical queries turns into a p99 spike on the database. Chapter 16.
Why "newer" is harder to decide than it looks, which matters for versioned sets and last-write-wins. Chapter 26.