Local vs distributed caching

Every pod keeps its own in-memory cache - fast, but they drift apart. Update the database, then watch local caches go stale while a shared Redis stays consistent. Invalidation, made visible.

9 min read

An in-process cache is the fastest thing in your system: no serialisation, no network, just a pointer in the same memory. It is also the thing most likely to serve a user data that stopped being true four minutes ago. Once an app runs as more than one pod, where the cache lives decides how wrong it can be.

Five pods, five different answers, all confidently wrong

Give each pod its own in-memory cache and nothing is shared, which is exactly why it is fast. It is also why a write to the database reaches none of them. Each pod keeps serving its own copy until its own TTL happens to expire.

DBDatabase holds v2 of user:42
4 stale copies
Pod 1
v1 stale
Pod 2
v1 stale
Pod 3
v1 stale
Pod 4
v1 stale
Every pod cached v1. The database has since moved to v2, and no pod knows.
4
Every pod cached v1, then the database moved on. Clear one pod and the rest stay wrong.

It gets worse than “stale for five minutes”. The pods filled their caches at different times, so their TTLs expire at different times, and for several minutes the answer a user sees depends on which pod the load balancer picked. Refreshing the page can flip the value back and forth.

The speed comes from not coordinating, and so does the staleness. They are the same property.

Move the cache out of the process and the drift goes away

A distributed cache like Redis gives the whole fleet one copy to agree on. An invalidation happens once and everyone sees it, because there is only one place for it to happen.

DBDatabase holds v2 of user:42
1 stale copy
Redis, shared by every pod
v1 stale
Pod 1
uses Redis
Pod 2
uses Redis
Pod 3
uses Redis
Pod 4
uses Redis
Redis cached v1. The database has since moved to v2, so the one shared copy is stale.
4
One copy to agree on: a single clear fixes every pod at once.

It is not free. A local read is tens of nanoseconds; a Redis read is a network round trip of hundreds of microseconds. That is still far faster than the database, but about four orders of magnitude slower than local memory. And every pod now depends on one more thing being up: a cache everyone agrees on is a cache everyone fails with.

Real systems use both, and choose their staleness window

The answer is almost never one or the other. Most systems end up with a small local cache with a very short TTL in front of the shared one, which keeps most of the local speed and bounds the disagreement to a window you picked.

Local cache in each pod
about 50 ns · TTL 2 s
61%
Shared Redis
about 300 µs · one truth
36%
Database
about 5 ms · durable
3%
Average read: 274 µs (Redis alone: 700 µs). Redis sees 39% of reads. Two users can disagree for at most 2 seconds.
2 s
A two-second local TTL takes most hot reads off Redis, and the staleness you accept is exactly the TTL.

Yes, that reintroduces staleness, deliberately and with a known bound. “Stale for up to two seconds” is a documented trade-off; “stale until some pod's five-minute TTL expires” is a bug. The local layer also soaks up hot keys before they ever reach Redis, which is often the real reason to add it.

So the decision is not “local or distributed”. It is “how long may two users disagree, and what is that worth?” Seconds of disagreement on a product description cost nothing. Seconds on an account balance become a support ticket. Pick the window first and the topology follows.

The short version

  • Per-pod caches are the fastest and drift apart, because nothing tells them about writes.
  • A shared cache gives one copy and one invalidation, for a network hop and a new dependency.
  • Layer a short-TTL local cache over the shared one to keep the speed with bounded staleness.
  • Choose how long users may disagree first; the cache design follows from that.