KnowSys

Caching, Properly

Follow one user profile, user:42, as an app stops asking the database for it every time. You'll see how much a small cache buys, what it has to forget, and the three ways it hurts you: a popular copy expires, a copy goes stale, or a fleet of servers hold different copies.

⏱ 50 min read◆ BeginnerAssumes: a terminal and Python; hash tables, a key-value store such as Redis and basic SQL help
Start reading

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.

Run 1,000,000 skewed requests through an LRU cache holding 1%, 5% and 20% of the keys
python
Python
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")
output
C++
 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 ms

Look 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 database10,000 reads/s × 0.101,000 /s
Hit rate 95%10,000 × 0.05500 /s
Hit rate 99%10,000 × 0.01100 /s
Hit rate drops from 99% to 90% during an incident1,000 / 10010× the load
Going from 90% to 95% halves the database load2×

?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 BY that 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 pathp50p99
In-process dict0.17 µs0.25 µs
Redis GET + json.loads, loopback65 µs153 µs
Postgres primary-key lookup, Unix socket64 µs143 µs
Postgres GROUP BY over the whole table141 ms278 ms
Redis 7.0.15 and PostgreSQL 16 in a shared 4-CPU Linux container, redis-py 8.1, psycopg 3. The table has 100,000 rows. 20,000 random reads per row (200 for the 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:

Cache-aside: user:42 is read, cached, then changed
App serverruns the pageCachesmall and fast: ~0.2 msDatabasesource of truth: ~5 msuser:7a.pnguser:42avatar a.pnguser:7a.pnguser:42a.png · TTL 60 sGET user:42
Step 1. A page view for 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.
1 / 7

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?

PatternOn a read missOn a writeGood forWatch out for
Cache-aside (look-aside)The application reads the database, then sets the cacheThe application writes the database, then deletes the keyMost application caches in front of a databaseRaces between a slow reader and a writer (section 9)
Read-throughThe cache library calls a loader, a function you supply that fetches the value from the databaseUsually paired with one of the write patternsIn-process caches like LoadingCache in Caffeine, a Java caching libraryLoader errors and timeouts now live inside the cache
Write-throughAs read-throughWrites go to the cache, which writes the database synchronouslyData that's read right after it's writtenEvery write pays both latencies; the cache fills with data nobody reads
Write-behind (write-back)As read-throughThe cache acknowledges, then writes the database laterAbsorbing write bursts, countersThe cache holds the only copy for a while. A crash loses acknowledged writes
Write-aroundAs cache-asideWrites skip the cache entirely; entries age out by TTLWrite-heavy data that's rarely re-readReaders 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?

Flowchart of a write-back cache. Reads and writes that miss locate a cache block, write its old data to lower memory first if it is dirty, and load the new data. Writes then change only the cache block and mark it dirty
Write-back in a CPU cache, where the name comes from. A write changes only the cached block and marks it dirty, and the memory below hears about it later, when that block is evicted (the 'is it dirty?' branch). Until then the cache holds the only up-to-date copy, which is exactly the risk of write-behind.Image: Flin00, CC0, via Wikimedia Commons

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.

A four-slot cache filled by A, B, C and D at times 0 to 3. E at time 4 replaces A. D is used again at time 5. F at time 6 replaces B
LRU on a four-slot cache. The number next to each key is when it was last used. A to D fill the cache, E evicts A, the least recently used, and touching D again at time 5 moves it to the front. When F arrives, B is now the oldest, so B goes.Image: Advaitjavadekar, CC BY-SA 4.0, via Wikimedia Commons

?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.

012345ringA · 2B · 0C · 1D · 3E · 0F · 1hand
A ring of six cached entries with a usage counter after each name. The hand starts at A.

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.

src/backend/storage/buffer/freelist.c
postgres/postgres @ REL_16_4 ↗
C
	/* 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:

A scan pushes the hot keys out of an LRU cache
Cache: 4 slotsoldest on the left, newest on the rightEvictedgone from the cacheuser:3colduser:9colduser:7hotuser:42hotuser:100scan, onceuser:101scan, onceuser:102scan, onceuser:103scan, once
Step 1. The cache holds four profiles. 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.
1 / 6

The loop is worse than the scan, because nothing ever recovers. Try to predict it:

Predict before you read on

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.

caffeine/src/main/java/com/github/benmanes/caffeine/cache/FrequencySketch.java
ben-manes/caffeine @ v3.1.8 ↗
java
  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:

DetailWhat it does
Sixteen 4-bit counters per longA 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 blockThe 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 × maximumAfter 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 twoOne 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:

A scan key meets the gate in a W-TinyLFU cache
WindowLRU · 1% of the cacheAdmission gatecandidate vs victimProbation20% of the main regionProtected80% of the main regionuser:98count 1user:9count 2user:7count 9user:42count 15user:99count 1
Step 1. A small cache. 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.
1 / 5

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:

StepWhat happens
AccessA read misses on key K. The loader fetches it, and the cache stores it
SketchThe access increments K's counters. The sketch sees every access, hit or miss
WindowK goes into the window. If K is part of a burst, it can get hits here
GateWhen the window overflows, its oldest entry, the candidate, meets the victim from probation. The sketch estimates both counts
OutcomeThe higher estimate stays. If the candidate loses, it's dropped and the main region doesn't change
PromotionAn 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:

caffeine/src/main/java/com/github/benmanes/caffeine/cache/BoundedLocalCache.java
ben-manes/caffeine @ v3.1.8 ↗
java
  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:

PolicyIdeaUsed in
LRUEvict the least recently usedMost libraries' default; the baseline
CLOCKLRU approximation with a reference bit or counter and a sweeping handPostgres buffers, OS page replacement
2Q / SLRUNew entries go to a probation area; a second hit promotes themPostgres 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 dataZFS. 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 admissionCaffeine (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 evictionsNewer; 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.

TraceLRULFUS3-FIFOW-TinyLFUOPT (ceiling)
Zipf 0.9, stable34.2%44.3%45.1%45.1%53.3%
Zipf + a 2,000-key scan every 10,000 reads33.2%44.2%45.1%44.9%53.3%
Hot set replaced halfway34.2%25.5%44.9%44.9%53.3%
Loop over 1,100 keys0.0%0.0%76.1%79.3%90.7%
Zipf 0.9, stable, cache of 5,00051.5%59.3%60.0%60.5%69.6%
Simulated in Python 3.14. Hit rates exclude the scan requests themselves, which can never hit. One seed; the policies are deterministic given the trace.

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.

HabitWhy
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 TTLA 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 expiryPast 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:

A hot key expires and every reader misses at once
Readersfour of the many app threadsCacheholds user:42Databasethe same query runs for each readerreader 1reader 2reader 3reader 4user:42TTL 1 s leftuser:42set four timesGET user:42
Step 1. Many readers want 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.
1 / 5

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:

singleflight/singleflight.go
golang/sync @ v0.8.0 ↗
Go
// 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.

A memcache lease during a stampede
Client AClient BmemcachedDatabaseget Kget Kwait and retrySELECTset K (token)get K
Step 1. Client A misses. memcached returns a lease token with the miss: A is now the one allowed to fill K.
1 / 6

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):

Output
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.

StrategyBackend calls in 10 sMost at onceReader p999 (median of 3)
Naive cache-aside2566452 ms
singleflight, per process32852 ms
Redis lock (SET NX), losers poll every 10 ms6158 ms
XFetch, β = 110–113–436 ms
Stale-while-revalidate, one refresher5–6124 ms
Redis 7.0.15 and redis-py 8.1 in a shared 4-CPU Linux container. The key was pre-filled, so these are expiry stampedes, not a cold start. Call counts were identical across three runs; tail latencies varied by 2–3× between runs on the shared container.

And as a decision table:

TechniqueCoordinationReaders wait?Cost
Coalescing (singleflight, LoadingCache)In-processYes, for the one fetchOne fetch per process per expiry
Lock or lease in the cacheSharedYes, and they pollA lock timeout to tune; the holder can die
XFetchNoneRarelyA few duplicate fetches; store Δ and expiry with the value
Stale-while-revalidateA lock for the refresherNoServes 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):

How cache-aside leaves a stale value behind
ReaderCacheDatabaseWriterGET K → missSELECT → v1UPDATE → v2DELETE KSET K = v1v1 until TTL
Step 1. The reader misses on K, which is not cached right now.
1 / 6

?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.

ApproachHow it worksWhere you see it
LeasesA miss issues a token; a delete invalidates it; a set with an invalid token is refusedmemcache 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 valuesStore 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 logA daemon follows the database's replication log, its ordered record of every committed change, and issues deletes after each commit, in commit orderFacebook'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

ApproachStalenessCostFails how
Short TTL only (seconds)Up to one TTLNothing to build; more missesWrong 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 delayA publish per write, a subscription per serverA server that misses a message stays stale until the TTL
Server-assisted tracking (Redis CLIENT TRACKING, Redis 6+)One message delayServer memory per tracked key and clientAs above, plus forced invalidations when the tracking table fills
Versioned keys (user:42:v17)None for readers who know the versionAn extra lookup to learn the current versionThe 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]:

Output
[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.

Five servers placed at positions 74, 139, 220, 310 and 360 on a circle. An object hashing to 111 is assigned to the server at 139
Consistent hashing puts servers and keys on the same circle. A key belongs to the first server clockwise from its hash, so the object at 111 goes to the server at 139. If that server dies, only the keys between 74 and 139 move, to the server at 220, and every other key stays where it was.Image: WikiLinuz, CC BY-SA 4.0, via Wikimedia Commons

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):

Shell
# 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 key

11.2Rules that hold up

  1. 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).
  2. Size the database for the misses. A cold or failed cache sends the whole read load to it (section 2).
  3. Jitter every TTL on keys that are written together.
  4. Guard every hot key against a stampede with coalescing, a lease, or early refresh (section 8).
  5. Use a scan-resistant policy if batch jobs share the cache with real traffic (sections 4 and 5).
  6. Give every local cache an invalidation path and a flush on disconnect (section 10).

11.3What you trade for what

You getYou payWhen the bill arrives
Reads served in microseconds instead of millisecondsData that can be stale for up to a TTLWhen a user sees an old profile
A database sized for the misses, not the readsThe database's capacity now depends on the hit rateOn 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 stepWhen you write it yourself instead of using a library
A lock or lease against stampedesA timeout to tune, and a holder that can dieWhen the winner crashes mid-fetch
Fast in-process copiesMany copies that must be kept coherentAs different servers showing different values

11.4Symptom, cause, fix

SymptomLikely causeFix
Database load spikes on a regular beatHot keys expiring together, or one key stampedingJittered TTLs, coalescing, XFetch or stale-while-revalidate
Database overloaded after a cache restart or deployCold cache; nothing to coalesceWarm from a peer, ramp traffic, rate-limit misses at the source
Hit rate falls after a nightly jobA scan evicting the hot set under LRUA 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 awayThe stale-set race, and no TTLAdd a TTL now; then leases, versioned sets or log-driven invalidation
Different servers show different valuesLocal caches with no invalidationShort TTL, broadcast invalidation, or Redis client tracking
One cache node at 100% CPU, others idleA hot keyLocal cache for that key, or replicate it under several names
Lots of lookups for IDs that don't existNo negative caching; possibly an enumeration attackCache "not found" with a short TTL
Hit rate fine, latency not improvedCaching something that was already cheapCache the expensive computation, or move the cache in process

12Summary

  1. 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%.
  2. 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.
  3. 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.
  4. Delete on write, don't update. Deletes are idempotent; out-of-order sets leave old values behind.
  5. LRU is simple and fragile. It's defeated by loops and hurt by scans, and a concurrent LRU writes on every read.
  6. Frequency needs aging. Plain LFU collapsed from 44% to 25.5% when the hot set changed; periodic halving fixes it.
  7. 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.
  8. Every TTL is a staleness budget and a safety net. Keep one even when you invalidate, because invalidations get lost.
  9. 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.
  10. 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.
  11. 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 NX lock, 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 SELECT and its SET. 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

check yourself
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?›
  1. 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 Δ.

Operating Systems: Three Easy Pieces, chapter 22

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.

Nishtala et al., Scaling Memcache at Facebook (NSDI 2013)

Look-aside caching at scale: leases, Gutter, mcsqueal, cold-cluster warm-up and cross-region consistency. PDF.

Einziger, Friedman and Manes, TinyLFU: A Highly Efficient Cache Admission Policy

The W-TinyLFU design and its evaluation on real traces. arXiv 1512.00727, later in ACM Transactions on Storage.

Vattani, Chierichetti and Lowenstein, Optimal Probabilistic Cache Stampede Prevention (VLDB 2015)

The proof behind XFetch, and why the exponential gap is optimal. PDF.

Yang et al., FIFO Queues Are All You Need for Cache Eviction (SOSP 2023)

S3-FIFO, one-hit wonders, and 6,594 production traces. PDF.

Caffeine: FrequencySketch.java and BoundedLocalCache.java

The best-documented production cache policy you can read. Start at admit() and FrequencySketch, then the hill climber. v3.1.8.

Redis: client-side caching reference

Tracking, broadcasting, the invalidation table, and the two races every near-cache must handle. redis.io.

Megiddo and Modha, ARC (FAST 2003)

The adaptive policy with ghost lists that most later work compares against. PDF.

RFC 5861: HTTP Cache-Control Extensions for Stale Content

stale-while-revalidate and stale-if-error, the HTTP names for serving old data on purpose. rfc-editor.org.

Redis Internals

Sampled LRU and the logarithmic LFU counter, expiry, and why a volatile-* policy can refuse writes. Chapter 22.

Memory Hierarchy & Cache Coherence

The same problems in hardware: cache lines, eviction and MESI keeping many copies coherent. Chapter 02.

Contention, Queueing & Tail Latency

Why a stampede's burst of identical queries turns into a p99 spike on the database. Chapter 16.

Time, Clocks & Ordering

Why "newer" is harder to decide than it looks, which matters for versioned sets and last-write-wins. Chapter 26.