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