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.
Design a multi-region caching strategy for a global application that requires sub-50ms read latency worldwide but strong consistency for user profile writes. Compare active-active replicated caches with conflict resolution, a global master for writes with local read caches, and CRDT-based approaches. Recommend an approach and justify trade-offs for latency, consistency, and operational complexity.
Sample Answer
Direct answer
There is no design that gives you both sub-50ms global reads and strong consistency for every write; you choose where the compromise sits. For most globally-distributed services the right default is regional read caches with a single-writer region per key (or per user) plus asynchronous replication, escalating to active-active with conflict resolution only for the specific keys that genuinely need multi-region writes.
Structured elaboration
- Global master with local read caches: writes always go to one authoritative region (chosen per key, e.g., the user's home region); that region's cache is updated synchronously, and other regions' caches are populated by asynchronous replication or event-driven invalidation. Reads in the writer's own region are strongly consistent; reads elsewhere are eventually consistent with a bounded lag (typically under a few seconds for a well-tuned replication pipeline). This is the simplest model to reason about and is the right default unless a specific workload proves it insufficient.
- Active-active with conflict resolution: every region can accept writes to the same key; conflicts are resolved with last-write-wins (using synchronized clocks or hybrid logical clocks), version vectors, or application-specific merge logic. This removes the single-writer bottleneck but pushes real complexity into conflict resolution and makes the failure modes harder to reason about; reserve it for keys where availability during a regional partition matters more than avoiding conflicts (e.g., a shopping cart, where a merge is acceptable).
- CRDT-based approaches: conflict-free replicated data types (structures with a mathematically well-defined merge function, such as a grow-only counter or a last-writer-wins register) let every region write locally and merge automatically without coordination. They are the cleanest fit when the data shape is naturally a CRDT (counters, sets, simple registers) but do not generalize to arbitrary application state.
- Choosing the read-latency-vs-consistency point per requirement: a common pattern is to segment data by freshness requirement. Authentication and permission state gets synchronous, single-region-authoritative reads (correctness matters more than latency). Profile pictures, preferences, and feed content get regional caches with eventual consistency (latency matters more, and staleness of a few seconds is invisible to users).
Worked example
For a service needing sub-50ms p95 reads worldwide with a single global-master design and 150ms average inter-region round-trip time: a user in Singapore reading data whose write-authoritative region is US-East would need a cross-region hop on every uncached read, blowing the 50ms budget by 3x. The fix is not "make replication faster" but "make the region the user reads from also the region their writes update" (their home region is the writer), with cross-region only used for failover. Under active-active with a 5-second bounded staleness target, replication lag has to be monitored per-region-pair and alerted when it exceeds roughly half the staleness budget (2.5s), leaving headroom for the merge/apply step.
Trade-offs and pitfalls
The most common design mistake is picking one consistency model for the whole service instead of per data class; forcing strong consistency on data that does not need it (e.g., a "last seen" timestamp) burns your entire latency budget for no user-visible benefit. Active-active without a clear conflict-resolution policy invites silent data loss (the losing write of a last-write-wins resolution simply disappears); make sure product owners have explicitly signed off on what "conflict" means for each data type before choosing this model. CRDTs remove coordination cost but restrict your data model; retrofitting CRDTs onto an existing schema is usually more work than designing around single-writer-per-key from the start.
Describe a multi-level caching architecture for a web service: L1 in-process cache (per instance), L2 shared cache (Redis cluster), and a CDN in front of static assets. Explain read and write flows, benefits for latency and throughput, and the primary consistency and invalidation challenges for each layer.
Sample Answer
Direct answer
Layer the caches by how expensive a miss at that layer is and how often the content changes: content delivery network (CDN) edge for static/shared assets closest to the user, a regional shared cache (e.g., Redis) for personalized-but-cacheable data, and an application in-process cache for the hottest few keys where even a network hop to the regional cache is too slow, with each layer's time-to-live (TTL) and invalidation strategy sized to that layer's job.
Structured elaboration
- CDN/edge layer: caches fully static or long-TTL content (images, JS/CSS bundles, and cacheable HTML fragments for anonymous/public views) at points of presence near the user. This absorbs the largest fraction of read volume for almost no cost per request and should hold anything that does not vary per-user.
- Regional distributed cache (Redis/Memcached): serves personalized or frequently-changing data shared across all app instances in a region (product details, prices, computed recommendations). Requests per second (RPS) here are much lower than at the edge because the edge already absorbed the static traffic, but this layer must handle write-driven invalidation correctly since content changes.
- Application in-process (L1) cache: for the small set of extremely hot keys (a handful of top-selling products, a config flag read on every request), an in-process cache avoids even the network round-trip to the regional cache. It is the fastest layer and also the hardest to keep coherent, because every app instance has its own copy; use a short TTL or an event-driven invalidation signal (pub/sub) rather than relying on the regional cache alone.
- Database-side / materialized views: for expensive aggregate queries, a materialized view or a query-result cache sits closest to the database, reducing load on the primary store even when the higher layers miss.
- Invalidation strategy per layer: static assets at the edge use content-hashed URLs so "invalidation" is really just a new URL (no purge needed); the regional cache uses event-driven invalidation on write (an order/price update publishes an invalidation for that product's key); the in-process L1 layer uses a short TTL (seconds) plus a lightweight pub/sub signal, because coordinating a purge across every instance is expensive and slow.
- Sizing to a concrete target: pick numbers for the design (RPS, latency budget, SKU/product count) and work backward: at, say, 100,000 RPS globally with a 90 percent edge-cacheable static-asset ratio, only about 10,000 RPS reach the regional application tier, which then needs to comfortably serve that load with sub-50ms p95 latency from cache.
Worked example
For a product catalog with 1,000,000 stock-keeping units (SKUs) and 50,000 RPS globally: static images and category pages are edge-cached (roughly 35,000 RPS absorbed at essentially zero backend cost). The remaining ~15,000 RPS of personalized/price-sensitive reads hit the regional Redis cache; with a realistic 95 percent hit ratio there, the origin database sees roughly 750 RPS, which is a design a mid-sized read replica tier handles comfortably. Price changes (a write path) publish an invalidation event per SKU, propagated to the regional cache and any in-process L1 caches holding that SKU within roughly 100 to 500 ms, well inside a "price must update within seconds" product requirement.
Trade-offs and pitfalls
Adding more layers adds more places for staleness to hide; a change that only invalidates the regional cache and forgets the in-process L1 layer will show correct data to some app instances and stale data to others, which is a confusing bug class to debug. Do not put personalization-sensitive content at the CDN edge unless you are using edge compute (e.g., edge functions) that can vary the response per user; naively caching personalized HTML at a shared edge node leaks one user's data to another. Every additional layer is also an additional operational surface (its own metrics, its own failure mode, its own on-call runbook); do not add a layer unless the sizing math shows the layer above it cannot meet the latency or load target alone.
Describe common cache eviction policies (LRU, LFU, FIFO, and TTL-based expiry). For each policy explain: how it decides what to evict, what access patterns it suits best, and give a short concrete example scenario (e.g., session caching, hot-working set, analytics counters).
Sample Answer
Direct answer
Least Recently Used (LRU) evicts by recency, Least Frequently Used (LFU) evicts by access count, First-In-First-Out (FIFO) evicts by insertion order regardless of access pattern, and time-to-live (TTL) based expiry evicts by a fixed age regardless of how popular the item is; each fits a different assumption about what predicts future value.
Structured elaboration
- LRU: evicts the item that has gone longest without being accessed. Fits workloads with strong temporal locality (a session that is hot right after login, then cold); weak spot is a one-time scan flushing genuinely popular items.
- LFU: evicts the item with the fewest accesses (often with decay). Fits a stable popularity distribution (a small set of perennially popular catalog items); weak spot is that brand-new items start at zero and can be evicted before they get a chance to prove their popularity.
- FIFO: evicts the oldest-inserted item regardless of how often it has been accessed since. It is the simplest and cheapest to implement (no access-order bookkeeping needed on every read), and fits workloads where recency of insertion genuinely correlates with relevance (e.g., a queue of recent events), but performs poorly for general-purpose caching where popular older items get evicted just because they were inserted first.
- TTL-based expiry: evicts (expires) an item after a fixed duration regardless of access pattern; it is really answering a different question ("how stale is too stale") rather than "what should we keep when we're full," and is often combined with one of the other policies as a ceiling on staleness rather than used as the sole eviction mechanism.
- Choosing for a workload: session caching (LRU: recency-driven), hot-working-set analytics (LFU or a hybrid: stable popularity), analytics counters with a natural expiry (TTL: staleness is the real constraint, not memory pressure).
Worked example
A cache of user session objects with occasional heavy churn (many users logging in and out in bursts) fits LRU well: sessions are hot right after login and cold afterward, and LRU naturally keeps the currently-active sessions warm through the churn. LFU would be a poor fit here because it would take time to "learn" a new burst of sessions are popular, during which it might evict genuinely active sessions in favor of older sessions that accumulated more historical hits.
Trade-offs and pitfalls
FIFO is cheap but usually the wrong default for a general-purpose cache; it is easy to reach for because it is simple to implement, but it ignores the one signal (access pattern) that actually predicts value. Do not conflate TTL (a staleness ceiling) with an eviction policy chosen for memory pressure; a cache can need both simultaneously for different reasons.
You must choose between Redis and Memcached to implement a session store for a web app. List trade-offs and recommend one choice. Consider persistence, data types, replication/HA, memory efficiency, eviction semantics, and operational features such as monitoring and backup.
Sample Answer
Direct answer
Choose Redis when you need richer data structures, persistence, or replication/high-availability (HA) built in; choose Memcached when you need the simplest possible pure key-value cache with the lowest per-operation overhead and multi-threaded read scaling out of the box.
Structured elaboration
- Data types: Redis supports strings, hashes, lists, sets, sorted sets, and more, useful when the cache itself needs to do more than store opaque blobs (e.g., a sorted set for a leaderboard). Memcached only stores simple key-value byte strings.
- Persistence: Redis can persist to disk (RDB snapshots, append-only file, AOF) so data can survive a restart; Memcached is purely in-memory with no persistence, so a restart is always a full cache flush.
- Replication / high availability: Redis has built-in replication and clustering (Sentinel for failover, Cluster for sharding); Memcached has no native replication, relying on the client or an external layer for HA.
- Memory efficiency: Memcached's simpler data model generally has lower per-key memory overhead for pure key-value use cases; Redis's richer data structures and features carry some additional overhead.
- Eviction semantics: both support least-recently-used (LRU) style eviction, but Redis offers more configurable policies (
allkeys-lru,volatile-ttl,allkeys-lfu, and others) versus Memcached's simpler slab-based LRU. - Operational features: Redis has a larger ecosystem for observability, Lua scripting for atomic multi-step operations, and pub/sub; Memcached's multi-threaded architecture can give it an edge on raw throughput for simple get/set at very high concurrency on a single node.
Worked example
For a session store: if sessions are pure key-value blobs, Memcached is a perfectly reasonable, simpler choice. If sessions need persistence across a restart (avoiding logging every user out simultaneously) or replication for HA, Redis with AOF persistence and Sentinel-managed failover is the safer choice, at the cost of a slightly more complex operational footprint.
Trade-offs and pitfalls
Choosing Redis "because it can do more" for a workload that is genuinely simple key-value adds operational surface area (persistence tuning, replication topology) without benefit; match the tool to the actual requirement. Memcached's lack of native replication means a node loss is a hard cache-miss event for everything that node held, with no automatic failover; that must be an explicit, accepted trade-off, not an oversight.
How would you design caching for large binary or JSON objects larger than 1MB that need to be served with low latency? Discuss the trade-offs of your approach, including memory-fragmentation considerations, and how to balance latency versus cost.
Sample Answer
Direct answer
For objects larger than about 1MB (megabyte), caching the object's bytes directly in a general-purpose in-memory cache is usually the wrong default; instead, offload the bytes to object storage and cache only a reference (and small, frequently-needed metadata), reserving direct in-memory caching for cases where the object is both large AND accessed with very low latency requirements that object storage cannot meet.
Structured elaboration
- Chunking: for objects that are read partially (a video seek, a large document's specific section), splitting into smaller chunks lets you cache only the actively-accessed portions rather than the whole object, and lets a partial cache hit still provide value.
- Compression: reduces both the memory footprint per cached object and network transfer time, at the cost of CPU for compress/decompress; worth it when memory or network bandwidth is the tighter constraint relative to available CPU.
- Offloading to object storage with cached references: store the actual bytes in a system built for large-object storage (durable, cheap per gigabyte) and cache only a reference (a URL or storage key) plus small metadata (size, content type, a short-lived signed access token if needed); this keeps the fast in-memory cache's precious capacity for many small, hot items rather than a few large ones crowding it out.
- Memory fragmentation considerations: storing objects of widely varying sizes in the same cache (a mix of small config values and occasional large blobs) can cause memory fragmentation in some cache implementations, reducing effective usable capacity below the raw configured limit; segregating large objects into a separate cache instance or tier, sized and tuned differently, avoids this.
- Balancing latency versus cost: caching large object bytes directly in memory gives the lowest possible latency but at high memory cost per item (crowding out many smaller, possibly more valuable cached items); offloading to object storage with a content delivery network (CDN) in front trades a small amount of latency (still fast, just not "already in local memory" fast) for much lower cost per byte and no fragmentation risk to the primary cache.
Worked example
A service serving large generated reports (5 to 50 MB each): rather than caching the full report bytes in the same Redis instance used for small, frequently-accessed session and configuration data, generated reports are written to object storage with a content-addressed key, and only that key (plus small metadata) is cached in Redis; clients fetch the actual report bytes from object storage (behind a CDN for repeat access) using the cached reference, keeping Redis's memory dedicated to the many small, hot items it is well-suited for.
Trade-offs and pitfalls
Caching large object bytes directly "because it's simple" can quietly degrade the whole cache's effectiveness for everything else sharing that cache instance, by consuming a disproportionate share of memory and contributing to fragmentation; segregate large and small objects into different tiers by default rather than only after noticing a problem. Offloading to object storage adds a small amount of latency and a second system to operate (object storage plus, often, a CDN in front of it); this is the right trade for genuinely large objects, but do not apply it reflexively to moderately-sized objects where direct in-memory caching remains the simpler, faster choice.
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.