Distributed cache
Design a cache cluster in front of a database for a read-heavy service, where a handful of keys take most reads. The hard part is not the misses but the hot keys.
- 1A hash ring spreads keys, not requests. With a skewed workload, the node that holds the hottest key runs far above the mean.
- 2A small L1 in each app server takes most hot reads off the network; it costs stale reads up to its TTL.
- 3Detect hot keys with a fixed-size sketch, then copy them to several nodes; each copy costs memory and one more fill.
- 4Cache-aside with delete on write and a lease per key keeps the cache close to the database through writes and failures.
Design a distributed cache in front of a database, where a handful of keys take most of the traffic.
- functional
- get(key), and invalidate on every database write.
- reads
- 1 million a second at peak; 1% writes.
- skew
- The 10 hottest keys take about half of all reads.
- latency
- p99 under 2 ms for a hit.
- freshness
- A write shows within 1 s.
- database
- Must stay under its read limit, also after a node failure.
I will shard the cache with a hash ring, then treat hot keys apart: an L1 in each app server, and copies on several nodes.
| question | answer assumed | it decides |
|---|---|---|
| How many keys, how large? | 10⁹ keys, 1 KB values | Memory, node count |
| How skewed are reads? | Zipf 1.25: top key 23% | Hot-key handling |
| How stale may a read be? | Up to 1 s | L1 and its TTL |
| Are hot keys also written often? | No: writes spread evenly | Copies, not split counters |
| Who owns invalidation? | The service that writes the row | Delete on write |
| What may the database take? | About 57k key reads a second | Required hit ratio |
| Managed or self-run cache? | Self-run Redis, client-side routing | Ring in the client |
I ask how skewed the reads are and how stale a read may be. Those two answers decide whether I need an L1 and hot-key copies.
- reads
- 1.00M a second at peak; ÷ 3 ≈ 333k average.
- writes
- 1% of requests ≈ 10k a second at peak.
- hottest key
- 23% × 1.00M ≈ 232k a second, on one node.
- data
- 10⁹ × 1 KB = 1 TB in the database.
- cache
- 16% of keys (the lab size) = 160 GB; × 1.5 for overhead ≈ 240 GB.
- nodes
- 240 GB / 32 GB ≈ 8 nodes, plus 1 replica each.
- misses
- (1 - 94.2%) × 1.00M ≈ 58k database reads a second.
- at nodes
- With the L1: (1 - 68%) × 1.00M ≈ 319k a second, ≈ 40k a node.
Hit ratio and L1 share come from the lab simulation at the same skew. The miss rate sits near the database limit, so the hit ratio is a requirement.
About 240 GB of cache on 8 nodes. The hottest key alone is 232k reads a second, more than one node serves without batching.
| call | result | on failure |
|---|---|---|
| get(key, loader) | value | Node down: call the loader |
| invalidate(key) | ok | Retry from a queue |
| GET key | value or nil | timeout 5 ms |
| SET key:lease token NX PX | OK or nil | Wait, then GET |
| SET key value EX | OK | Ignore: next read fills |
| DEL key key#0..7 key:lease | count | Retry; TTL bounds it |
- loader reads the database. The library calls it at most once per key per process.
- invalidate is idempotent: a second DEL deletes nothing.
- Every cache call has a short timeout. A slow cache is worse than a miss.
The cache has no API of its own beyond get and invalidate. Every other rule lives in the client library that the service uses.
- products
- One row per key, with a version that grows on each write.
- product:<id>
- The whole row, so one GET answers one read. The client sends the key's lease to the same node, so one script can check both.
- #0 to #7
- Copies of a hot key. Each name hashes to its own ring point, so most land on different nodes.
- L1
- An LRU of 100 entries in each of 16 app servers. It holds the hottest keys, because they are read most.
- sketch
- 64 counters per app server, reset each second.
The database row is the truth. The cache holds it under one key on its ring owner, plus 8 copies when it is hot; each app server keeps a small L1.
Step 1: Read, L1 hit
- Look the key up in this process's LRU of 100 entries.
- A hit costs no network call.
If it fails
The entry can be up to 1 s older than the database, because a write does not reach other servers' L1.
An app server answers from its L1, then from the cache node that owns the key, then from Postgres. Writes commit, then delete every copy.
| tool | capability | what it gives this design | also used for |
|---|---|---|---|
| Client library | Consistent hash ring, 100 tokens a node | Any client finds a key's node; a node loss moves only its keys. | Sharded databases, load balancers |
| App server | In-process LRU with a TTL | Hot reads cost no network call. | Config, feature flags |
| App server | Space-Saving sketch | Finds hot keys in fixed memory. | Top-k queries, abuse detection |
| Redis | GET, SET EX, DEL of several keys | The shared cache; one DEL removes a key and its copies. | Sessions, rate limits |
| Redis | SET NX PX with a token | A lease: one filler per key, and no stale fill after a write. | Locks, idempotency keys |
| Redis | maxmemory with an LRU or LFU policy | The node evicts cold keys itself when memory is full. | Any bounded cache |
| Redis | Pipelining | Many GETs in one round trip; one node serves far more a second. | Bulk loads |
| Postgres | Primary-key read | The fill on a miss. | Every lookup |
The ring places keys, the L1 and copies spread hot reads, and leases with delete on write keep the cache close to the database.
| mitigation | status | busiest node |
|---|---|---|
| None: one home per key | Not approved | 3.35× mean |
| L1 in each app server | 1 s stale is fine | 1.21× mean |
| Copies of hot keys on the ring | Approved | 1.40× mean |
| Both | Approved | 1.21× mean |
| Read replicas of the hot node | Stale replica reads are fine | Spreads one node's reads |
| Split a hot key into sub-keys | Hot writes, such as counters | A read sums every part |
- Find hot keys by sampling requests or with a sketch. Counting every key exactly costs memory per key.
- Space-Saving keeps 64 counters. A key above 1/64 of the reads is always in the table.
- The L1 needs no detection: an LRU of 100 entries keeps the hottest keys by itself.
Lab: 600,000 requests, 50,000 keys, Zipf 1.25, 16 app servers, 8 nodes. Seeded; it repeats exactly.
on each read of key at app server a:
a.sketch.add(key) // 64 counters, Space-Saving1
every 1 s at app server a:
a.hot = keys whose count > 1% of a's reads2
a.sketch = new sketch
route(key):
IF key is in a.hot:
RETURN key + "#" + random(0 .. 7) // one of 8 copies3
RETURN key // its one home on the ring4
read(key):
v = a.L1.get(key) // 100 entries, 1 s TTL5
IF v: RETURN v
node = ring.owner(route(key))
v = node.GET(...) OR fill from the database
a.L1.put(key, v)- 1A new key takes the smallest counter and inherits its count. Hot keys stay; cold keys churn through the last slots.
- 2Each app server decides alone. With round-robin load balancing, every server sees the same skew. The lab found 12 hot keys.
- 3Each copy hashes to its own ring point, so two copies can share a node. Colder keys also stay uneven. The lab shows 1.4×, not 1.0×.
- 4Cold keys keep one copy, so memory stays close to one entry per key.
- 5The L1 answered 68% of reads in the lab.
Tested source Go: detect · Go: route · Go: read · Go: Space-Saving, reused from the B10 lab
// detect closes one window for an app server: every key whose Space-Saving count is above
// HotShare of the reads it saw becomes hot for the next window. The sketch holds Counters
// keys, so the memory stays fixed however many keys pass through.
func detect(a *app, cfg Config) {
hot := map[string]bool{}
for _, it := range a.ss.Top(cfg.Counters) {
if float64(it.Count) > cfg.HotShare*float64(a.n) {
hot[it.Key] = true
}
}
a.hot, a.n = hot, 0
a.ss = sketches.NewSpaceSaving(cfg.Counters)
}
// cacheKey picks where a read goes. A key the app server marked hot has Copies copies,
// key#0 to key#(Copies-1); each read takes one at random, so the copies hash to different
// nodes and share the load. Every other key has one home on the ring.
func cacheKey(a *app, cfg Config, key string, r *rand.Rand) string {
if cfg.Replicate && a.hot[key] {
return fmt.Sprintf("%s#%d", key, r.IntN(cfg.Copies))
}
return key
}
res.Reads++
reads[k]++
bReads++
a.ss.Add(key)
a.n++
if cfg.L1 {
if v, ok := a.l1.Get(key); ok { // 1. this app server's memory
res.L1Hits++
bHits++
if v < version[k] {
res.StaleReads++
}
continue
}
}
ck := cacheKey(a, cfg, key, r)
n := nodes[ring.Owner(ck)] // 2. the cache node that owns the key
res.NodeLoad[n]++
bNode[n]++
v, ok := caches[n].Get(ck)
if ok {
res.NodeHits++
bHits++
if v < version[k] {
res.StaleReads++
}
} else { // 3. the database, then fill the cache
res.DBReads++
bl.DBReads++
v = version[k]
caches[n].Put(ck, v, time.Hour)
}
if cfg.L1 {
a.l1.Put(key, v, cfg.L1TTL)
}
// Add counts one occurrence of key. A key in the table gets +1. A new key takes a free counter,
// or replaces the key with the smallest count and inherits that count + 1. The inherited part is
// its error bound.
func (s *SpaceSaving) Add(key string) {
s.n++
if i, ok := s.h.at[key]; ok {
s.h.items[i].Count++
heap.Fix(&s.h, i)
return
}
if s.h.Len() < s.cap {
heap.Push(&s.h, &Item{Key: key, Count: 1})
return
}
low := s.h.items[0]
delete(s.h.at, low.Key)
s.h.at[key] = 0
*low = Item{Key: key, Count: low.Count + 1, Err: low.Count}
heap.Fix(&s.h, 0)
}
With no mitigation the busiest node ran at 3.35× the mean. An L1 brought it to 1.21×, and copies to 1.40×.
| method | status | stale window |
|---|---|---|
| Set the new value on write | Not approved | Until TTL: two writers can race |
| TTL only, no invalidation | Stale for TTL is fine | Up to the TTL |
| Delete on write | Rare fill race is fine | A fill that read before the write |
| Delete on write, lease on fill | Approved | None in the shared cache |
| An L1 in front | 1 s stale is fine | Up to the L1 TTL |
| Change stream deletes keys | Many writers | The stream's lag |
- Lab: 0 stale reads without an L1 and 9,842 with one, out of 593,976 reads.
- The L1 is never invalidated by other servers' writes. Its TTL is the freshness promise.
- Keys that must never be stale skip the L1.
write(key, row):
database.update(row) // commit first1
DEL key, key#0 .. key#7, key:lease // every copy, and the lease2
fill(key) after a miss:
IF SET key:lease token NX PX 2000: // one filler per key3
v = database.read(key)
SET key = v ONLY IF the lease is still token4
ELSE: wait 10 ms, then GET again- 1Delete before the commit, and a reader can fill the old row back in the gap.
- 2The writer does not know which keys other servers call hot, so it deletes all 8 copy names.
- 3A lease is a lock with a token. The B4 lab measured 1 database query for 100 readers of one expired key.
- 4The write path deleted the lease, so a filler that read the old row cannot store it.
Tested source Go: delete every copy · Go: B4 write, then delete · Go: B4 fill under a lock
// write is cache-aside with delete on write: after the database commits, delete the key and
// every copy it may have on the ring. The next read fills from the database.
func write(caches []*caching.LRU[string, int64], ring *sharding.Ring, nodes map[string]int, cfg Config, key string) {
drop(caches[nodes[ring.Owner(key)]], key)
if cfg.Replicate {
for c := range cfg.Copies {
ck := fmt.Sprintf("%s#%d", key, c)
drop(caches[nodes[ring.Owner(ck)]], ck)
}
}
}
// UpdatePrice is the write path: change the row, commit, and only then delete the cache key and
// its lease. The next reader fills the cache from the committed row.
func (c *Cache) UpdatePrice(ctx context.Context, db *pgxpool.Pool, id, price int64) (version int64, err error) {
err = db.QueryRow(ctx, sqls["update_price"], id, price).Scan(&version) // commits on its own
if errors.Is(err, pgx.ErrNoRows) {
return 0, ErrNotFound
}
if err != nil {
return 0, fmt.Errorf("update product %d: %w", id, err)
}
return version, c.Invalidate(ctx, id)
}
// Invalidate deletes the cached value and any open lease, in one command.
func (c *Cache) Invalidate(ctx context.Context, id int64) error {
if err := c.rdb.Del(ctx, Key(id), leaseKey(id)).Err(); err != nil {
return fmt.Errorf("invalidate %d: %w", id, err)
}
return nil
}
func (c *Cache) fillUnderLock(ctx context.Context, id int64) (Product, error) {
token := rand.Text()
for {
won, err := c.rdb.SetNX(ctx, lockKey(id), token, c.LockTTL).Result() // SET NX PX
if err != nil {
return Product{}, fmt.Errorf("take fill lock %d: %w", id, err)
}
if won {
return c.fillAndUnlock(ctx, id, token)
}
// Another reader is filling. Wait for the value, or for the lock to vanish.
for {
if err := sleep(ctx, c.PollEvery); err != nil {
return Product{}, err
}
if p, hit, err := c.Peek(ctx, id); hit || err != nil {
return p, err
}
n, err := c.rdb.Exists(ctx, lockKey(id)).Result()
if err != nil {
return Product{}, fmt.Errorf("check fill lock %d: %w", id, err)
}
if n == 0 {
break // the holder died without filling: try to take the lock
}
}
}
}
func (c *Cache) fillAndUnlock(ctx context.Context, id int64, token string) (p Product, err error) {
defer func() {
if uerr := unlock.Run(ctx, c.rdb, []string{lockKey(id)}, token).Err(); uerr != nil {
err = errors.Join(err, fmt.Errorf("release fill lock %d: %w", id, uerr))
}
}()
// Check again: the previous holder may have filled the key just before we took the lock.
if p, hit, err := c.Peek(ctx, id); hit || err != nil {
return p, err
}
return c.fill(ctx, id)
}
I commit to the database first, then delete every copy of the key. A lease per key stops a slow filler from writing back an old row.
| cache holds | hit ratio | database reads |
|---|---|---|
| 2% of keys | 85.5% | 8,625/s |
| 4% of keys | 89.0% | 6,509/s |
| 8% of keys | 91.9% | 4,802/s |
| 16% of keys | 94.2% | 3,463/s |
| 32% of keys | 95.7% | 2,557/s |
| 64% of keys | 96.0% | 2,380/s |
| eviction | status | keeps |
|---|---|---|
| LRU over all keys | Approved | Recently read keys |
| LFU over all keys | Stable hot set | Often read keys |
| Only keys with a TTL | Mixed with data that must stay | Keys without a TTL |
| No eviction | Not approved | Writes fail when memory is full |
Lab rows: 8 nodes, a 10 s run. The first read of each key always misses, which caps the hit ratio near 96%.
a cache node fails:
clients drop it from the ring // its keys move to the next tokens1
each moved key misses once2
one filler per key: singleflight in the process,
a lock across processes3
a node joins:
it takes about 1/(n+1) of the keys, all cold
add it while traffic is low, or copy the hottest keys to it first4- 1Only that node's keys move. Hash mod N would move almost every key and empty the whole cache.
- 2The level after the failure stays a little higher, because 7 nodes hold fewer entries than 8.
- 3A hot key that moves gets thousands of misses at once. One filler, and the rest wait.
- 4A cold node takes its share of misses all at once. A cold start in the lab sent 3× the warm level to the database.
Tested source Go: B4 singleflight · Go: B4 fill under a lock
// Do runs fn for key, or waits for the call already running for key. shared is true when the
// result came from another caller's call.
func (g *Group[V]) Do(key string, fn func() (V, error)) (v V, err error, shared bool) {
g.mu.Lock()
if c, ok := g.calls[key]; ok {
g.mu.Unlock()
<-c.done // join the call in flight
return c.val, c.err, true
}
c := &call[V]{done: make(chan struct{})}
if g.calls == nil {
g.calls = map[string]*call[V]{}
}
g.calls[key] = c
g.mu.Unlock()
c.val, c.err = fn() // only this caller queries
g.mu.Lock()
delete(g.calls, key)
g.mu.Unlock()
close(c.done)
return c.val, c.err, false
}
func (c *Cache) fillUnderLock(ctx context.Context, id int64) (Product, error) {
token := rand.Text()
for {
won, err := c.rdb.SetNX(ctx, lockKey(id), token, c.LockTTL).Result() // SET NX PX
if err != nil {
return Product{}, fmt.Errorf("take fill lock %d: %w", id, err)
}
if won {
return c.fillAndUnlock(ctx, id, token)
}
// Another reader is filling. Wait for the value, or for the lock to vanish.
for {
if err := sleep(ctx, c.PollEvery); err != nil {
return Product{}, err
}
if p, hit, err := c.Peek(ctx, id); hit || err != nil {
return p, err
}
n, err := c.rdb.Exists(ctx, lockKey(id)).Result()
if err != nil {
return Product{}, fmt.Errorf("check fill lock %d: %w", id, err)
}
if n == 0 {
break // the holder died without filling: try to take the lock
}
}
}
}
func (c *Cache) fillAndUnlock(ctx context.Context, id int64, token string) (p Product, err error) {
defer func() {
if uerr := unlock.Run(ctx, c.rdb, []string{lockKey(id)}, token).Err(); uerr != nil {
err = errors.Join(err, fmt.Errorf("release fill lock %d: %w", id, uerr))
}
}()
// Check again: the previous holder may have filled the key just before we took the lock.
if p, hit, err := c.Peek(ctx, id); hit || err != nil {
return p, err
}
return c.fill(ctx, id)
}
Hit ratio grows slower with each doubling of memory. When a node dies, its keys miss once on their new nodes: database reads rose from 303 to 418 per 100 ms in the lab.
600,000 requests · 50,000 keys · Zipf 1.25 · 16 app servers · 8 cache nodes · L1 100 entries, 1 s
bars: requests that reached cache nodes 0 to 7, out of 593,976 reads · dashed line: the mean · scale fixed to the run with no mitigation
- busiest node
- 3.35× mean
- served by L1
- 0%
- hit ratio
- 94.2%
- database reads
- 3,463/s
- stale reads
- 0
An L1 takes the hottest reads off the network, and copies spread what is left. Each one costs something: stale reads for the L1, memory and fills for the copies.
| pattern | GETs a second | busiest node | batch p99 |
|---|---|---|---|
| Spread over 10,000 keys | 1,539,327 | 30% | 2.7 ms |
| One hot key | 808,651 | 100% | 1.9 ms |
| Hot key, a copy per node | 1,507,210 | 25% | 3.0 ms |
- Each Redis node runs one event loop for commands, so one key can use only one node's loop.
- Copies give the hot key the whole cluster again.
- Spread keys are a little uneven: the ring gave one node 30% of the reads, not 25%.
Redis 8, 4 server processes on one 16-thread laptop, 64 clients, 16 GETs per round trip, 1.5 s a run. Rates vary between runs; the test checks only the shares and the order.
On 4 real Redis nodes, one hot key reached 809k GETs a second, all on one node. Spread keys reached 1.54M, and 4 copies of the hot key 1.51M.
| event | result | why it is safe | saved by |
|---|---|---|---|
| One key gets a quarter of all reads | Its node runs at several times the mean. | The L1 serves most of those reads; copies spread the rest. | L1, copies |
| A cache node dies | Its keys miss once on their new nodes. | The ring moves only that node's keys; leases make one filler per key. | Ring, leases |
| A hot key expires | Many readers miss at once. | One reader holds the lease and fills; the others wait 10 ms. | Lease |
| A write and a fill race | The fill read the old row. | The write deleted the lease, so the fill is refused. | Lease |
| The DEL after a write fails | The shared cache keeps the old value. | Retry from a queue; the TTL bounds the window if all retries fail. | Retry queue |
| A write changes a key in L1 | Other servers serve the old value. | For at most the L1 TTL, which the requirements allow. | L1 TTL |
| The whole cache is cold | Most reads go to the database. | Warm the hottest keys before traffic, or let traffic in slowly. | Warm-up |
| Memory is full | New keys need space. | LRU or LFU eviction drops cold keys; the database has every row. | maxmemory |
| step | add | it handles | move up when you see |
|---|---|---|---|
| 1 | One Redis, cache-aside, delete on write, a lease per key. | About 36k single GETs a second (B4 lab); one machine's memory. | Memory full with a low hit ratio, or the node near its rate. |
| 2 | A ring of nodes in the client, with a replica each. | Memory and throughput grow with nodes; a node loss moves only its keys. | One node far above the mean: a hot key. |
| 3 | An L1 in each app server, with a short TTL. | 68% of reads with no network call in the lab. | The hot node is still the busiest, or the L1 TTL is too stale for some keys. |
| 4 | Hot-key detection and copies. | Busiest node near 1.2× the mean. | Users in other regions, or a region must survive alone. |
| 5 | A cache cluster per region, invalidated by the database's change stream. | Local latency everywhere; each region warms its own keys. | Top of the ladder. |
Demand is from panel C. Single-GET and Postgres rates are the B4 lab lows; batched rates are from panel L.
I start with one Redis and cache-aside. I add nodes on a ring when memory or throughput runs out, and hot-key handling when one node runs hot.
| question | answer | go deeper |
|---|---|---|
| Why not Redis Cluster's slots instead of a ring? | Slots work the same way: 16,384 fixed slots, and a node owns a set of them. Clients learn the map and follow redirects when a slot moves. | X2 |
| The hot key is a counter that is written all the time. Now what? | Split it into N sub-keys and add to a random one. A read sums all N, so reads cost more. | X2 |
| How do you find hot keys in production? | Sample a share of requests, or keep a count-min sketch with a small heap in each client. Report the top keys each window. | B10 |
| What stops a stampede when a hot key expires? | A lease or singleflight lets one reader fill; others wait. Refreshing early, before expiry, removes the miss. | B4 |
| Why not write through the cache? | The cache then needs the database's transactions and failure handling. Cache-aside keeps the database the only writer of truth. | B4 |
| How do you invalidate across regions? | Read the database's change stream in each region and delete the keys there. The lag of the stream is the stale window. | X1 |
0 of 9 known
The ring spreads keys evenly. Why does one cache node still run hot?
The L1 served 68% of reads, yet the hit ratio stayed at 94.2%. Why?
What does the L1 cost?
Why does copying hot keys add database reads?
How does an app server find hot keys without a counter per key?
Why delete the cache key on a write instead of setting the new value?
A cache node dies. What does the database see?
Cache memory doubles from 16% to 32% of the keys. How much does the hit ratio gain?
On real Redis, what limits one hot key?
- skew
- Zipf 1.25: top key 23% of reads, top 10 55%.
- busiest node
- 3.35× mean; L1 1.21×; copies 1.40×.
- L1
- 100 entries per server served 68% of reads.
- hit ratio
- 16% of keys cached: 94.2%. 4%: 89.0%.
- Redis
- One hot key: 809k GETs a second on one node; 4 copies: 1.51M.
- latency
- L1 hit about 150 ns; Redis GET about 0.2 ms (B4 lab).
- failure
- 1 of 8 nodes lost: database reads 303 → 418 per 100 ms.
Simulation runs are seeded and repeat exactly. Redis rates: Redis 8 on a 16-thread laptop shared with other jobs; they vary between runs.