skip to content

You run four instances of the same application against one database, each with its own Hibernate second-level cache. What stale-data problems appear, and how do invalidation-based and replicated or distributed cache topologies differ in handling them?

level: seniorimportance: must knowfreq 40%

answer

  1. N JVMs, N caches, one database
  2. Invalidation = send a tombstone, reload on miss
  3. Replication = ship the value, full copy per node
  4. Distribution = N owners, remote reads
  5. Async propagation = real stale window; TTL bounds it

basics

~20 s

A commit on one node leaves the other three holding the old rows until they are told. Invalidation topologies broadcast "forget this key" so other nodes reload from the database; replicated or distributed topologies ship the new value instead. Asynchronous messaging leaves a stale window either way.

solid answer

~50 s

With independent per-node caches, only the writing node's cache is corrected on commit; the others serve pre-commit state indefinitely. That is unacceptable for anything but truly immutable data, so a clustered provider is used. **Invalidation** sends a small "key X is dead" message after commit. Cheap, scales with write rate, and each node reloads on its next miss — so the cost lands on the database. **Replication** ships the new state to every node: reads stay local after a write, but every node pays memory for the whole dataset and every write costs a full-payload broadcast. **Distribution** keeps N copies spread across the cluster: bounded memory, but many reads become a remote fetch. Whichever you pick, propagation is usually asynchronous, so there is a window where node B answers with the old value. Bound it with a short time-to-live, use synchronous invalidation for the few entities that need it, and accept that a strictly fresh read must bypass the cache.

go deeper

for a junior

Know that each application instance has its own second-level cache, so a write on one node does not fix the others unless the cache provider is clustered.

for a middle

Contrast invalidation with replication in terms of message size, memory, and where the reload cost lands.

for a senior

Reason about the asynchronous stale window, TTL as a safety net, which entities may be cached at all, and how restarts and lost messages behave.

for a principal

Frame it as a consistency and capacity decision: pick the topology from dataset size, write rate and staleness tolerance, define which reads may never be served from cache, and plan refill behaviour for deploys and failures.

## The failure mode The second-level cache lives inside a JVM. Scale to four instances and you have four caches over one database. Node A commits an update: its own after-commit callback fixes its cache, but nodes B, C and D never saw the write and keep serving the old row. Nothing expires it, so without help the staleness is unbounded. This is the single most common production surprise when a service that was cached happily on one instance is scaled out. ## Topologies **Invalidation.** After commit, the writing node broadcasts a message naming the key (or region) that is now invalid. Each receiving node removes its local copy; the next request for that key misses and reloads from the database. Message size is tiny and independent of entity size, so write-heavy workloads stay cheap on the network. The cost is database load: every node reloads the row separately, so a hot row that changes often is fetched N times per change. This is the usual default for entity regions in a clustered setup. **Replication.** Every node holds a full copy of the data and every write is propagated as state, not as a tombstone. Reads are always local, even right after someone else's write, which is attractive for small, read-mostly reference data. The costs are memory (the dataset must fit in every node) and network (each write ships a payload to every node), so it scales poorly with dataset size and write rate. **Distribution.** Each key is kept on a fixed number of owner nodes (say two out of eight). Total memory is bounded and the cluster can grow, but a read on a non-owner is a remote call, which is far slower than a local hit and can be slower than a well-indexed database read. It suits large datasets that must be cached but cannot fit per node. ## The stale window Propagation is normally asynchronous for throughput. Between node A's commit and node B applying the invalidation, node B answers reads with the pre-commit value. The window is usually milliseconds but is not bounded by anything in the application. Synchronous propagation shrinks it at the price of adding cluster round trips to your commit path, and it still cannot make the cache transactional with the database unless you run a fully transactional cache mode, which is expensive and rarely justified. Practical ways to bound the exposure: - Configure a short time-to-live so any missed or lost invalidation self-heals. - Cache only entities whose staleness is tolerable; keep balances, inventory counts, prices at checkout time, and permission decisions out. - For the rare read that must be fresh, bypass the cache explicitly (a query with cache mode set to ignore/refresh, or a pessimistic lock which forces a database read). ## Related cluster hazards - **Split brain and dropped messages.** If the invalidation bus partitions, nodes drift apart silently. A time-to-live is your only automatic recovery. - **Rolling deploys with changed mappings.** Two versions of an entity class sharing one distributed cache is a deserialization or semantics bug waiting to happen; version the region name or run separate regions per app version. - **Cache clocks.** Cached query results are validated against per-table update timestamps; with independent clocks across nodes those comparisons get fuzzy, which is one more reason cached query results are the first thing to switch off when a cluster misbehaves. - **Restart storms.** A node starting cold pulls its whole working set from the database. Rolling restarts of a big cluster serialise those refills onto one database. ## How to answer State the problem (N independent caches, one database), name the three topologies with their cost profile, be explicit that asynchronous propagation means eventual consistency with a real stale window, and finish with the mitigations plus the honest boundary: if a read must be correct to the transaction, it must not come from a shared cache.

  • Which topology would you choose for a small country-code reference table read on every request?
    Replication. The dataset is tiny, effectively read-only, and the value of always having a local hit outweighs the negligible cost of propagating rare writes. Invalidation would give the same freshness but force every node to re-query the database after each change, and distribution would add pointless remote reads for data that fits everywhere.
  • How do you make one particular read immune to cross-node staleness?
    Take it out of the cache path: set the cache retrieve mode to bypass or refresh for that query, or acquire a pessimistic lock, which forces Hibernate to read the row from the database with a locking SELECT. Both trade the performance benefit for a guarantee, which is the right trade only for the few reads that need it.
  • An invalidation message is lost because of a brief network partition. What limits the damage?
    A time-to-live or maximum idle time on the cache region: the stale entry eventually expires and the next read reloads it. Without expiry the entry can stay wrong until the process restarts, so a TTL is standard practice even when invalidation is believed to be reliable.

saying these in an interview costs you the question

  • Assuming the second-level cache is automatically shared across instances just because it is enabled
  • Claiming a clustered cache gives the same consistency as the database
  • Recommending full replication for large or write-heavy datasets
  • Ignoring the asynchronous propagation window when reasoning about correctness
  • Having no expiry configured because "invalidation handles it"

context