Caching Strategies and Distributed Caching Questions
Using caches to reduce latency and load: cache-aside, read-through, write-through, and write-behind patterns, TTLs, eviction policies, and distributed caches such as Redis or Memcached. Covers cache invalidation, stampede and thundering-herd protection, and the consistency tradeoffs of caching. Focuses on where and how to cache across tiers.
Discuss how to handle cache serialization and deserialization safely and efficiently. Consider versioning serialized formats, schema evolution, backward compatibility, and lazy migration strategies during rolling upgrades.
Sample Answer
Direct answer
Cache serialization needs explicit schema versioning so that a deployment change (a field added, removed, or retyped) does not produce a runtime error or silent data corruption when new code reads an old cached format, or old code reads a new one, during a rolling upgrade.
Structured elaboration
- Versioning the serialized format: embed a schema version alongside the serialized payload (either in the value itself, or in the cache key, as covered by the key-versioning pattern) so a reader can detect which shape it is dealing with.
- Schema evolution: prefer additive, backward-compatible changes (new optional fields default sensibly if absent) over breaking changes (renaming or retyping an existing field) wherever possible, since additive changes let old and new code coexist safely without any special handling.
- Backward compatibility during rolling upgrades: during a deploy, old and new application code run simultaneously for some window; if new code writes a new schema version and old code cannot deserialize it, old-code instances will error or silently misbehave on every cache read for that key until the deploy completes, unless the new schema is designed to be backward-readable, or the version is used to explicitly route old code away from new-format entries.
- Lazy migration strategies: rather than migrating every cached entry to a new format immediately, let entries migrate lazily: on a cache miss (or an explicit read-and-rewrite), the new format is written; old-format entries simply age out via their normal time-to-live (TTL), avoiding a disruptive bulk-migration pass.
- Choosing a serialization format: a format with strong built-in support for optional/default fields (protocol buffers, for example) makes additive schema evolution far more natural than a format that requires exact structural matching to deserialize correctly.
Worked example
Adding a new optional field to a cached user-profile object: with a backward-compatible serialization format, old code (not yet aware of the new field) simply ignores it when reading a new-format entry, and new code supplies a sensible default when reading an old-format entry that lacks the field; no special versioning logic is needed at all, because the format itself tolerates the difference. Contrast with renaming an existing field, which is NOT safely backward-compatible under most formats and would require either a version-gated read path or accepting a brief window of cache misses/errors during the rolling deploy.
Trade-offs and pitfalls
Assuming a serialization format is "safe" without checking its specific behavior on missing or extra fields is a common way this bites during a rolling deploy; verify the format's actual compatibility guarantees, do not assume. Skipping explicit versioning "because most changes are additive" works until the first genuinely breaking change arrives unplanned; building the versioning mechanism in from the start costs little and avoids an emergency retrofit later.
You must size an in-memory cache cluster for a dataset: total data size 200GB, expected working set 20GB, replication factor 2 for high availability, and 20% extra headroom for fragmentation and metadata. Explain how you would calculate node count and per-node memory, consider shard overhead, and account for future growth.
Sample Answer
Direct answer
Size a cache cluster from the actual working set, not the total dataset, then add replication and fragmentation headroom on top; node count and per-node memory come out of that total, not the other way around.
Structured elaboration
- Working set versus total data size: the total dataset (200 GB, gigabytes, in a typical example) is not what needs to fit in cache; the working set (20 GB, the subset actually read repeatedly within a relevant time window) is the real sizing driver, since the whole point of caching is serving the repeatedly-accessed subset, not the entire dataset.
- Replication factor: for high availability, each byte of the working set is stored on more than one node (replication factor 2 doubles the effective memory requirement, since every item exists on a primary and at least one replica).
- Fragmentation and metadata headroom: real-world memory usage is always somewhat higher than the raw data size, due to per-key overhead (metadata, pointers) and memory fragmentation from the allocator; a 20 percent headroom is a reasonable starting estimate, refined with actual measurement once running.
- Calculating node count and per-node memory: total required memory equals working set size times replication factor times (1 plus headroom fraction); node count times per-node memory must exceed that total, with per-node memory chosen based on available instance types and a preference for more, smaller nodes (better fault isolation, smaller blast radius per node loss) versus fewer, larger nodes (simpler operations, less network overhead).
- Accounting for future growth: size with headroom for expected growth over a planning horizon (e.g., 6 to 12 months), not just current working-set size, since scaling a cache cluster up later involves the same resharding/migration considerations covered elsewhere in this topic.
Worked example
Working set 20 GB, replication factor 2, 20 percent headroom: required memory is 20 GB×2×1.2=48 GB. Choosing nodes with 8 GB of usable cache memory each (leaving room for the node's own operational overhead), that requires ⌈48/8⌉=6 nodes. If 6-month growth projections suggest the working set could grow to 30 GB, planning for 30×2×1.2=72 GB, or 9 nodes at the same per-node size, avoids a resharding event shortly after initial rollout.
Trade-offs and pitfalls
Sizing off total dataset size instead of working set drastically overestimates the memory needed for most read-heavy workloads with meaningful access skew (a small fraction of data accounting for most reads); always validate the working-set assumption with real access-pattern data where available, not just an estimate. Under-provisioning headroom for fragmentation and metadata is a common source of "sized correctly on paper, still evicting more than expected in production" surprises; treat the headroom percentage as a starting estimate to refine with real measurement, not a fixed constant to trust blindly.
At extreme scale, a single cache miss for a hot key can overload the origin. Propose a comprehensive defense-in-depth strategy to prevent stampedes: singleflight, background regeneration, early recompute, probabilistic TTLs, prewarmed hot key paths, and rate limiting. Explain how to orchestrate these across many app instances.
Sample Answer
Direct answer
At extreme scale you cannot rely on a single stampede defense. Combine request coalescing (only one request repopulates a hot key while others wait), background/proactive refresh before expiry, probabilistic early expiration, jittered time-to-live (TTL) values, and origin rate limiting as a last-resort backstop, coordinated so all application instances agree on who is allowed to refresh a given key at once.
Structured elaboration
- Request coalescing (singleflight): on a cache miss, the first request acquires a short-lived lock (e.g.,
SETNXin Redis) for that key and fetches from origin; concurrent requests for the same key either block briefly on a notification channel (pub/sub or polling with backoff) or serve a stale value if one exists. This bounds concurrent origin load per key to roughly one in-flight fetch, regardless of instance count, because the lock lives in the shared cache, not in any one process. - Probabilistic early expiration: instead of a hard expiry, each read close to TTL end recomputes with a small, increasing probability (a common formula is P(refresh)=e−β⋅(texpiry−tnow)/δ where δ is the time it took to compute the value and β tunes aggressiveness). This spreads refreshes across many requests instead of concentrating them at the exact expiry instant.
- Jittered TTLs: add randomized jitter (e.g., base TTL plus/minus 10 to 20 percent) so keys written around the same time do not all expire in the same millisecond, which is what turns an ordinary cache miss into a correlated stampede across thousands of keys at once (a cache avalanche).
- Background/prewarmed refresh: a scheduled worker (or the request that detects "close to expiry") refreshes hot keys proactively so the TTL rarely actually lapses for high-traffic keys; this trades a small amount of continuous background load for eliminating stampede risk on the hottest paths.
- Origin rate limiting as a backstop: even with the above, cap concurrent origin requests per key (or per origin endpoint) so a defense-in-depth failure degrades gracefully into serving stale data or a fast error instead of taking the origin down.
- Orchestrating across many app instances: the lock, the "who refreshes next" decision, and the notification of waiters must live in the shared cache (Redis) or a coordination service, not in in-process state, because coalescing only works if every instance agrees on a single winner per key.
Worked example
Say a hot key normally takes 200 ms to recompute and serves 5,000 requests per second (RPS) at peak. Without any defense, if it expires with no coalescing, the next ~1,000 requests in that 200 ms window (5,000 RPS times 0.2 s) would all miss and hit the origin simultaneously; a database that comfortably serves single-digit concurrent queries per second for that expensive query falls over. With coalescing, exactly 1 request recomputes and the other ~999 either wait ~200 ms for the notification or receive the previous (slightly stale) value immediately; origin load for that key stays at 1 concurrent request regardless of RPS.
Trade-offs and pitfalls
A lock that never expires on a crashed refresher permanently blocks that key; always set a lock TTL slightly longer than the expected recompute time, plus a fallback path that lets a waiter give up and fetch directly after a bounded wait. Serving stale-while-refreshing is a correctness trade-off, not a free win: it is the right default for read-heavy, staleness-tolerant data, and the wrong default for a low-latency-but-must-be-fresh field like an account balance. Jitter alone does not help an already-hot key that is legitimately read far more than others; that is a hot-key sharding problem, not a stampede problem, and needs a different fix (splitting the key, adding a replica-backed local cache).
How do you decide whether to introduce a cache for a given service endpoint? Describe the signals and measurements you would collect, the tests you would run (load, latency, profiling), and the criteria that justify adding an in-process cache, a shared cache (Redis), or a CDN. Include considerations for cost, operational complexity, and correctness.
Sample Answer
Direct answer
Decide whether to add a cache by measuring the actual read pattern (how often the same value is requested, how expensive it is to produce, and how much staleness is tolerable), not by defaulting to caching every endpoint; a cache with a low repeat-read rate or zero staleness tolerance is a cost with no real benefit.
Structured elaboration
- Signals to collect: request rate for the same key/query (does the same data actually get read repeatedly, or is nearly every read unique), the cost of producing the value (a fast, cheap lookup gains little from caching even if repeated), and the data's staleness tolerance (how quickly must a change be visible).
- Tests to run: a load test comparing latency and backend load with and without a proposed cache, and a profile of the actual query/computation to confirm it is genuinely a meaningful cost worth caching against.
- Criteria for choosing a cache tier: an in-process cache fits data that is cheap to duplicate per instance and does not need cross-instance consistency; a shared cache (Redis) fits data that benefits from being consistent across instances or too large to duplicate per instance; a content delivery network (CDN) fits public, non-personalized content that benefits from being close to users geographically.
- Cost: weigh the infrastructure and operational cost of adding a caching layer (a new dependency to monitor, secure, and keep available) against the actual load/latency benefit measured above; a marginal benefit may not justify the added complexity.
- Operational complexity: caching adds invalidation logic, a new failure mode (cache unavailable), and another thing to monitor; these costs are real even when the caching decision is otherwise sound, and should be weighed explicitly.
- Correctness: if the data's staleness tolerance is effectively zero (a value that must always reflect the absolute latest state, with no acceptable delay), caching adds risk without benefit, since any caching mechanism introduces at least a small window of potential staleness.
Worked example
An endpoint returning a real-time stock quote, requested uniquely per symbol per user with essentially no repeat reads within any meaningful window, and requiring zero staleness tolerance: this fails on both the "does the same value get read repeatedly" test and the "can staleness be tolerated" test, making it a poor caching candidate regardless of how expensive the underlying computation is. Contrast with a product description, read thousands of times per hour by different users for the same handful of popular items, changing rarely: this passes both tests clearly.
Trade-offs and pitfalls
Caching by default, without measuring the actual read-repetition rate, either wastes cache capacity on data that gets no benefit or, worse, introduces a staleness risk on data that could not tolerate it; always start from measurement, not habit. The decision is not binary per endpoint; the same service can have some data that benefits enormously from caching and other data (even on the same page) that should never be cached, and treating the whole endpoint uniformly misses that nuance.
Explain the purpose of caching in distributed systems. Define cache hit, miss, and hit ratio, and describe the typical benefits and tradeoffs of adding a cache to a service. Give concrete examples of workloads that benefit from caching (and why) and workloads where caching could be harmful.
Sample Answer
Direct answer
A cache is a fast, smaller copy of data kept close to where it is used, so repeated reads can be served without redoing expensive work (a database query, a network call, a heavy computation) every time. The core trade-off is speed and reduced load in exchange for the risk that the cached copy becomes stale relative to the real source of truth.
Structured elaboration
- What problems caching solves: latency (serving from memory or a nearby node is far faster than recomputing or fetching from a distant source), throughput/load (fewer requests reach the expensive backend), and cost (fewer database queries or external application programming interface (API) calls, which often cost money directly).
- Where caches typically live: client/browser (closest to the user, zero network cost on a hit), content delivery network (CDN) / edge (near the user but shared across many users), application in-memory (fast, but local to one process/instance), and a shared distributed cache like Redis or Memcached (shared across all instances in a region, one network hop away).
- Primary trade-offs: staleness (the cached copy may not reflect the latest write), added complexity (invalidation logic, cache-miss handling, monitoring another moving part), and memory cost (caches are not free storage).
- When caching helps: read-heavy workloads where the same data is requested repeatedly and can tolerate at least a little staleness, or where the underlying computation/fetch is expensive relative to a cache read.
- When caching is harmful: data that changes on every read (no repeated value to cache), workloads that are already write-heavy with low read repetition (cache churn without benefit), or correctness-critical data where any staleness is unacceptable and the added complexity of cache invalidation introduces more risk than the latency win is worth.
Worked example
An API endpoint that computes a dashboard aggregate over a large dataset in 800ms, requested 200 times per minute by the same handful of users, is a strong caching candidate: a 30-second time-to-live (TTL) cache serves nearly every request from cache after the first, cutting both latency (single-digit milliseconds instead of 800ms) and backend load by over 95 percent, for a staleness window most dashboard users will not notice. The same 30-second TTL applied to an account balance shown right after a deposit would be actively harmful, showing the user a stale, "wrong" number at the exact moment they are checking it.
Trade-offs and pitfalls
The most common mistake is caching by default rather than by evaluating whether the read pattern and staleness tolerance actually justify it; caching adds a second place data can be wrong (the cache disagreeing with the source of truth), and that failure mode does not exist at all if you never cache the data in the first place.
Unlock Full Question Bank
Get access to all Caching Strategies and Distributed Caching interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.