skip to content

Your service runs on several instances that each cache rows in memory; after one instance writes, why do the others keep serving the old value, and what makes them agree?

level: seniorimportance: should knowfreq 55%

answer

  1. one key, several separate entries
  2. invalidation stops at the process boundary
  3. the bug follows the load balancer
  4. broadcast narrows, expiry bounds
  5. sharing costs a round trip

basics

~20 s

Each process holds its own entry, so an invalidation on one instance is invisible to the rest. Making them agree needs a broadcast removal, a single shared out-of-process copy, or not caching mutable rows per instance at all.

solid answer

~50 s

An in-memory cache is per process: the same key in another instance is a different entry, and neither the database nor the write path reaches across. So the writer's invalidation is local and every other instance keeps its copy until something else clears it, which is why the value looks correct on some requests and wrong on others depending on routing. Three fixes exist. Broadcast the removal over a channel — reads stay in-process and fast, but a lost or reordered message is a silent stale entry. Move the entries to one shared store — agreement becomes structural, at the cost of a network round trip and serialisation on every hit, and identity is lost. Or cache only near-immutable data locally and read mutable rows live. Whichever you pick, keep a bounded maximum entry age so a lost invalidation self-heals.

go deeper

for a junior

Hold on to the core fact: an in-memory cache belongs to one process. Two instances with the same key hold two different copies, and removing one does not remove the other.

for a middle

Describe the two mechanisms that make them agree — broadcasting the removal, or keeping one copy in a shared store — and what each costs on the read path.

for a senior

Recognise the symptom from the field: correct on some requests, wrong on others, follows the load balancer. Then say how you bound it, with a maximum entry age behind whichever propagation you chose, and how you measure convergence.

for a principal

Decide per dataset rather than globally: which data may live in per-instance memory at all, whether a shared store's extra dependency is worth structural agreement, and what convergence time you are willing to promise.

## Why the instances disagree A cache that lives in a process's own memory is invisible to every other process. Two instances of the same service, running the same build, holding the same key, hold **two separate entries**. When instance A writes a row and correctly removes its own entry, instance B's entry is untouched: nothing in the database, and nothing in the write path, reaches into another process. B keeps serving its copy for as long as that copy is allowed to live. This produces the failure mode teams find hardest to reproduce: the value is correct on some requests and wrong on others, depending purely on which instance the load balancer picked. Refreshing the page "fixes" it, then unfixes it. Sticky routing does not solve it either — it only hides the divergence until a deployment, a scale-out, or a rebalance moves the user to another instance. ## Option one: broadcast the invalidation Keep the entries local and add a channel: when an instance commits a write, it publishes "this key changed" and every instance removes its own copy. - Reads stay in-process, so the fast path keeps its speed — this is the reason to choose it. - Delivery is typically best-effort. A message lost while an instance was restarting or partitioned is a stale entry that nothing will ever clear, and no counter reports it. - Ordering is not guaranteed. A message that arrives before the writer's commit is visible can cause a listener to remove the entry, immediately miss, reload the pre-image, and repopulate — the same race that the invalidate-before-and-after-commit ordering fights locally, now stretched across the network. - The window is at least the propagation delay, so the honest bound on staleness is the propagation delay plus whatever backstop age you configure. ## Option two: one shared copy Move the entries out of process into a store all instances read: now there is exactly one copy and agreement is structural rather than eventual. - Every hit becomes a network round trip plus deserialisation, so the cache saves the query but not the trip. For rows that were cheap to fetch, this can be a net loss. - The value must be reduced to bytes, which means you keep raw field values rather than a live object graph: identity is gone, and links to other objects have to be re-resolved. - The store becomes a dependency of every read, so its availability, its own maximum age, and its failure behaviour become part of yours. - It removes divergence between instances; it does **not** remove staleness. A write that bypasses the layer leaves the shared entry exactly as wrong, just wrong for everybody consistently. | Approach | Read cost | Instances agree | Main risk | |---|---|---|---| | Local entries, no coordination | in-process | no | unbounded divergence between instances | | Local entries plus broadcast removal | in-process | eventually | a lost or reordered message is silent | | One shared out-of-process store | network round trip | yes | added dependency; identity and links lost | | Cache nothing mutable per instance | full read each time | trivially | more load on the database | ## Option three: narrow what is cached The cheapest fix is often to stop caching the rows that change. Long-lived local entries are a good fit for near-immutable reference data, where minutes of divergence are unobservable, and a poor fit for rows a user edits and then immediately re-reads. Splitting the two — long-lived local entries for reference data, no cache or a short shared one for mutable rows — removes the class of bug rather than managing it. ## What to do about the messages you will lose Whatever you build, assume some invalidation will not arrive: treat the broadcast as an **optimisation that shortens the window**, and a bounded maximum entry age as the **guarantee that closes it**. Then the worst case is a number you can state — "no instance serves a value more than N seconds after the row changed" — instead of "until someone restarts the pod". Stamping entries with the row's version helps too: an instance that is handed an older stamp than the one it holds can refuse the overwrite, which kills the reordering race without any extra coordination. Finally, verify it rather than assuming it. Write a row, then read it back through every instance and record how long each took to converge. That measurement, repeated automatically, is what turns an invalidation design into a claim you can defend.

  • Why does sticky routing not fix instance-to-instance divergence?
    It only hides it. A user pinned to the instance that did the write sees fresh data, so the bug disappears in testing. A deployment, a scale event, or a rebalance moves that user to an instance holding an old copy, and the wrong value returns at the worst possible moment.
  • A broadcast removal arrives at another instance before the writer's commit is visible. What happens?
    The listener removes its entry, immediately misses, reloads the pre-image because the commit is not visible yet, and repopulates with the old value. It is the local invalidate-before-and-after-commit race stretched across the network, and it is why the message should be published after the commit and backed by a bounded entry age.
  • Does moving to one shared store remove staleness?
    No, only divergence. Every instance now agrees, but a write that bypasses the layer leaves that single entry just as wrong — consistently wrong for everybody instead of wrong for some. The bypass problem is orthogonal to where the entries live.
  • How would you verify the instances actually converge?
    Measure it. Write a marker row on a schedule, read it back through every instance, and record the time until all of them return the new value. Chart the maximum rather than the mean and alert when that tail crosses the staleness budget you have stated.

Every instance is a colleague with a photocopy. Correcting the original in the filing cabinet changes nobody's copy until you walk round and tell each of them.

saying these in an interview costs you the question

  • Assumes an invalidation propagates to other processes automatically
  • Thinks sticky sessions solve cross-instance divergence
  • Believes a shared store also fixes bypassing writes
  • Treats a broadcast channel as guaranteed delivery
  • Ignores the round trip a shared cache adds to every hit
  • Says a transaction boundary clears longer-lived entries