Interview Setup
Interview Prompt
Design a distributed in memory cache system (like Redis Cluster or Memcached) that supports get/set/delete with TTL, horizontal scaling across nodes, and high availability. Clients should see a single logical cache with automatic key distribution.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| What is the data size per key and the total working set? | Contrasting 100-byte session keys against 1 MB media objects directly determines memory sizing, slab allocation efficiency, and eviction pressure. |
| What is the read-to-write ratio and what are the consistency requirements? | A 1000:1 read ratio favors aggressive cache-aside strategies, whereas write-heavy workloads require write-through or proactive event-driven invalidation. |
| Can we tolerate transient cache loss on node failure? | Cache-aside architectures tolerate transient node outages by repopulating misses from the database, whereas session stores require replication to prevent dropped user sessions. |
| Do clients connect directly to nodes or route through an intermediate proxy? | Smart clients that maintain direct cluster slot mappings eliminate proxy network hops, whereas proxies like Twemproxy simplify client connection pooling at the cost of an extra network hop. |
Scope
In scope
- Consistent hashing with virtual nodes
- LRU and LFU eviction policies
- Cache stampede prevention
- Write-through, write-behind, write-around
- Hot key mitigation
- Node add/remove with minimal key movement
Out of scope (state explicitly)
- Persistent storage / AOF/RDB durability internals
- Full Redis command set (streams, pub/sub)
- Building a new consensus protocol for cache metadata
Functional Requirements
Start by clarifying whether the system requires basic string key-value semantics or structured collections, TTL expiration, eviction policies, and replication guarantees:
PUT(key, value, TTL): Store a key-value pair with an optional time-to-live duration.GET(key): Retrieve the value associated with a key, returning null if the key is absent or expired.DELETE(key): Explicitly remove a key-value pair from storage.- Support configurable eviction policies including LRU, LFU, FIFO, and Random.
- Support structured data types beyond raw strings, including hashes, lists, sets, and sorted sets.
- Support per-key TTL and expiration with active and passive cleanup.
- Support atomic primitives such as increment, decrement, and compare-and-swap.
- Distribute keys across multiple nodes for horizontal scalability and high availability.
Non-Functional Requirements
Interviewers prioritize sub-millisecond latency, hit ratio stability, and partition tolerance. Establish consistent hashing for shard placement and an eviction strategy early in the discussion:
- Ultra-Low Latency: Sub-millisecond p99 latency (< 1 ms) for in-memory read and write operations.
- High Throughput: 1M+ operations per second per node through connection pooling and pipelining.
- High Availability: Survive single-node and rack failures without total cache loss via automated replica failover.
- Horizontal Scalability: Dynamically add nodes with minimal key redistribution to scale memory capacity.
- Eventual Consistency: Bounded eventual consistency across replicas under best-effort cache semantics.
- Memory Efficiency: Maximize usable payload density per gigabyte of RAM using optimized memory allocators.
- Partition Behavior: Keep failures isolated and avoid cascading overload, while making stale or unavailable reads explicit during network partitions.
Capacity Estimations
| Metric | Calculation | Value |
|---|---|---|
| Total data cached | Given working set assumption | 10 TB |
| Avg key size | Given metadata assumption | ~100 bytes |
| Avg value size | Given payload assumption | ~1 MB |
| Total keys (approx) | 10 TB ÷ 1 MB per value | ~10M |
| Entry overhead (metadata) | Pointers, TTL expiry, allocator headers | 100 bytes per entry |
| Entries per node (64 GB RAM) | 64 GB ÷ ~1.0002 MB per entry | ~64K entries |
| Nodes needed (raw) | 10 TB ÷ 64 GB RAM per node | ~160 nodes before headroom and replication |
| Primary nodes with 70% usable RAM | 10 TB ÷ (64 GB x 0.70) | ~224 nodes |
| Total cluster nodes with replication factor 2 (1 primary + 1 replica) | 224 primary nodes x 2 (or 160 raw x 2) | ~320 nodes at raw capacity, or ~448 nodes with 70% usable RAM |
| Operations / sec (total peak) | Given peak throughput assumption | 100M |
| Operations / sec / node | 100M ops/sec ÷ node count | ~625K at 160 nodes, or ~446K at 224 primary nodes |
Architecture Diagram
In the room: default to cache-aside (lazy loading) unless your interviewer specifically requests write-through or write-behind semantics.
Walk your interviewer through the architecture following the client request path. Application servers route keys using the configured sharding scheme. Custom sharding can use consistent hashing, while Redis Cluster uses its deterministic hash slot mapping. Each shard replicates state asynchronously to replicas for high availability, allowing reads to scale across replicas while writes target the primary. When memory reaches provisioned capacity, an eviction policy such as LRU discards cold entries, while background timers expire entries based on TTL without requiring explicit delete calls.
Component Deep Dives
Walk through each component in the distributed cache hierarchy systematically, examining the in-process L1 tier, the client routing library, sharded data placement, in-memory allocators, and replication failover mechanics.
L1: In-Process Cache (Caffeine / Guava)
An in-process L1 cache absorbs extreme read spikes directly inside application heap memory, eliminating network latency before requests reach a distributed Redis cluster. Multi-level caching, sharding, and eviction form the foundational mechanisms that interviewers examine.
- Zero-Hop Latency: Embedded directly within each application server JVM or process memory, delivering sub-0.1ms latency without network hops.
- Memory Footprint: Approximately 100 MB per instance to ensure cache allocations do not exhaust JVM heap or trigger garbage collection pauses.
- Bounded TTL: Short durations of 10 to 30 seconds limit staleness because the L1 cache does not receive synchronized push invalidations from Redis.
- Eviction Strategy: Window TinyLFU (used by Caffeine) maintains near-optimal hit ratios and outperforms classic LRU for skewed access patterns.
- Use Cases: Highest-frequency keys including dynamic configuration flags, active user sessions, and viral product catalog entries.
- Architectural Trade-off: Each application server maintains an isolated L1 cache, meaning updates can leave different instances returning divergent values until the short TTL expires. This eventual consistency is acceptable for read-heavy catalog data, but unsuitable for linearizable data models.
Request flow: 1. Check L1 in-process cache. On hit, return value immediately (< 0.1ms) 2. On L1 miss, check L2 distributed cache (Redis): on hit, return value and backfill L1 (< 1ms) 3. On L2 miss, read from primary database: return value, populate L2, and backfill L1 (5-50ms)
Cache Client Library
Sub-millisecond read SLOs require an intelligent client library. Embedded directly within each application process, the client manages sharded routing, connection pooling, compression, and circuit breaking without proxy hops.
- Hash Slot Routing: For Redis Cluster, computes
CRC16(key) % 16384, looks up the target node in its cached slot map, and transmits requests directly to the responsible master shard without proxy overhead. Custom consistent hashing deployments instead use the configured hash ring. - Slot Map Cache: Stores the complete cluster slot-to-node topology in application memory, refreshing the table automatically upon receiving
MOVEDredirects or via periodic 60-second polling. - Connection Pooling: Maintains persistent TCP socket pools of 10 to 20 connections per cache node, avoiding recurring TCP handshake and TLS negotiation overhead.
- Pipelining: Batches multiple commands into a single network round trip, reducing per-command overhead and network round trips during batch lookups.
- Compression: Automatically compresses payloads larger than 1 KB with LZ4 prior to transmission, reducing memory footprint and network bandwidth transparently.
- Circuit Breaker: Isolates failing nodes when five consecutive health checks fail, halting outbound requests to that instance for 30 seconds and falling back to the primary database or stale L1 data to prevent cascading thread starvation.
- Timeouts and Retries: Enforces a strict 2ms request timeout with a single retry, ensuring a lagging cache node does not block client threads.
Consistent Hashing: Data Distribution
Consistent hashing organizes keys and nodes onto a circular identifier ring, ensuring that scaling the cluster up or down reallocates only a minimal fraction of the total keys.
Why Simple Modulo Hashing Fails
Using naive modulo hashing where node = hash(key) % N causes nearly every key to remap whenever the cluster size N changes. Adding or removing even a single node invalidates approximately (N - 1) / N of cached entries, triggering a catastrophic cache miss storm that overwhelms the primary database.
Consistent Hashing Architecture
Consistent hashing projects both keys and physical nodes onto a continuous 32-bit circular ring spanning 0 to 2^32 - 1. Each key is assigned to the first node encountered moving clockwise from its hashed position. When a new node joins, it only assumes responsibility for keys located between itself and its immediate counterclockwise predecessor, which limits key remapping to roughly 1/N of the total keyspace. Similarly, when a node is removed or fails, only its local keys migrate to its clockwise neighbor, leaving the remainder of the cluster undisturbed.
Virtual Nodes (Vnodes)
To prevent hash skew where uneven node distribution creates hot shards, each physical node is assigned 100 to 200 virtual node positions distributed evenly across the ring. Virtual nodes smooth out key density across the cluster, ensuring that when a physical host fails, its workload disperses across multiple surviving servers rather than concentrating onto a single neighbor.
Ring positions: vnode_A1: 1000 vnode_B1: 2500 vnode_A2: 5000 vnode_C1: 6000 vnode_B2: 7500 vnode_C2: 9000 Key "user:123" hashes to 3000 -> assigned to vnode_A2 (position 5000)
In-Memory Storage Engine and Layout
Each cache shard maintains an in-memory hash table providing average O(1) time complexity for lookups, inserts, and deletes, paired with specialized memory allocators to prevent fragmentation under continuous write churn.
Core Storage Engine
The primary storage layer operates like an optimized C++ hash table (such as std::unordered_map). High-throughput hash functions like MurmurHash3 or xxHash minimize collision probability across keys. Collisions are resolved through separate chaining with linked lists or open addressing with Robin Hood hashing. To prevent latency spikes during table resizing, progressive incremental rehashing migrates entries across two internal tables in small increments when the load factor exceeds 0.75.
Memory Layout Optimization
Continuous allocation and deallocation of variable-length cache values can cause severe external memory fragmentation, leading to premature out-of-memory errors despite free RAM. Systems employ specialized allocators to maintain high payload density:
- Slab Allocator (Memcached Approach): Pre-allocates fixed 1 MB memory pages partitioned into uniform chunk classes (such as 64B, 128B, 256B, 512B, 1KB, up to 1 MB), placing each entry into the smallest fitting chunk class to eliminate external fragmentation.
- Jemalloc (Redis Approach): Utilizes an advanced general-purpose memory allocator with multiple memory arenas and fine-grained size bins, minimizing fragmentation while dynamically adapting to variable-length payloads without rigid chunk boundaries.
Eviction Policies and Lifecycle Management
When memory reaches provisioned capacity limits, an eviction policy determines which keys to discard to accommodate incoming writes. Selecting the right policy directly governs hit ratio stability under skewed or shifting access distributions.
Least Recently Used (LRU)
LRU evicts the entry that has gone unaccessed for the longest duration, optimizing for strong temporal locality where recently referenced data is likely to be queried again.
- Data Structure: Combines an O(1) hash map with a doubly linked list tracking access order.
- Read Updates: On a successful
GET, the corresponding node is unlinked and spliced to the head of the doubly linked list. - Eviction Execution: When memory is exhausted, the tail node representing the stalest key is evicted in O(1) time.
Approximated LRU (Redis Approach)
Strict LRU requires 16 bytes of pointer metadata (prev and next) per entry, which consumes gigabytes of RAM across hundreds of millions of keys. Redis avoids this overhead through probabilistic sampling:
- Sampling Algorithm: When memory exceeds
maxmemory, the engine samples a configurable number of random keys (typically 5) and evicts the candidate with the longest idle time. - Memory Efficiency: Sampling eliminates explicit linked-list pointers while approximating true LRU behavior within a few percentage points of accuracy.
Least Frequently Used (LFU)
LFU evicts keys with the lowest cumulative access frequency, making it ideal for power-law workloads where 20% of catalog keys drive 80% of read traffic. Unlike LRU, a one-off sequential table scan will not flush established hot keys from memory. Redis implements LFU using an 8-bit Morris logarithmic counter paired with configurable exponential decay, allowing historical popularity to decrease gradually so older items do not occupy cache memory indefinitely.
TTL-Based Expiration
Keys with an explicit time-to-live are reclaimed through a complementary pair of passive and active expiration mechanisms:
- Passive Expiration: Evaluates key expiration during read lookups, immediately deleting the expired record and returning a cache miss.
- Active Expiration: A background timer periodically samples a random batch of volatile keys (such as 20 keys ten times per second) and purges expired entries to reclaim RAM before client queries occur.
Replication and Failover Architecture
Replication decouples read scaling from write throughput while protecting cache availability during physical hardware failures.
Cluster Replication Model
Each primary shard coordinates with one or more asynchronous read replicas. Writes are executed on the primary shard and acknowledged immediately to preserve sub-millisecond latency, after which modifications stream asynchronously to replicas. Reads target the primary by default to ensure read-your-writes consistency, while read replicas can absorb offloaded queries for workloads that tolerate replication lag.
- Heartbeat Failure Detection: Shards exchange periodic gossip heartbeats to monitor node health and detect network partitions.
- Automated Replica Promotion: When a primary fails to respond within the configured timeout, surviving peers conduct an election to promote an eligible replica with the highest replication offset to primary status.
- Consistency Trade-off: Because replication operates asynchronously, an acknowledged write can be lost if the primary crashes before streaming updates to its replicas. In caching architectures, this brief window of data loss is acceptable because missing entries are backfilled on demand from the database.
Cache Coherence and Stale Write Protection
Cache-aside invalidation can race with concurrent reads when a database update and a cache refill execute concurrently. To maintain coherence, the application commits the database modification before invalidating the cached entry, and incorporates version numbers or monotonic timestamps into cache values. When refilling the cache after a miss, the client verifies that the database version is newer than the currently cached version before writing. For enterprise architectures requiring durable invalidation propagation across distributed regions, a transactional outbox or Change Data Capture (CDC) pipeline publishes versioned invalidation events to Kafka, providing durable and replayable delivery across cache clusters. If network delivery is delayed, the configured TTL bounds normal cache lifetime, subject to the cache expiration and serving policy.
- Commit-First Invalidation: Commit the primary database transaction before issuing cache invalidation requests to prevent premature cache repopulation with old data.
- Versioned Refill Guard: Attach version numbers to cached payloads to reject refill attempts where a slow concurrent read could overwrite newer committed data.
- Durable Invalidation Streaming: Propagate invalidation notices through an asynchronous change stream or CDC pipeline to provide durable and replayable delivery across regional cache clusters.
- Bounded TTL Fallback: Enforce an upper bound on cache lifespan so uninvalidated entries eventually expire even if network partitions interrupt pub/sub invalidation events.
Cluster Architecture (Redis Cluster Hash Slots)
Redis Cluster partitions the keyspace across 16,384 deterministic hash slots, providing clients with predictable routing across independent master shards without requiring a centralized proxy layer. This deterministic slot model contrasts with consistent hashing rings used by Memcached and Twemproxy topologies.
- Deterministic Hash Slots: The cluster divides the keyspace into 16,384 discrete slots distributed evenly across primary nodes. Clients compute
CRC16(key) % 16384to identify the responsible slot. - Explicit Slot Allocation: Each primary master owns an assigned range of slots (for example, Node A manages slots 0 to 5,460, Node B manages 5,461 to 10,922, and Node C manages 10,923 to 16,383).
- Direct Client Routing: Smart client libraries cache the slot-to-node mapping locally, transmitting commands directly to the target primary socket without intermediary hops.
- MOVED Redirect Protocol: During cluster rebalancing or node migration, a contacted node that no longer owns the requested slot returns a
MOVED slot target_ip:portresponse, instructing the client to execute the command on the new node and update its internal slot table.
API Design
Core Cache Operations
The client interface exposes key-value primitives and atomic operations, establishing clean domain contracts for application developers:
export type CacheValue = string | Uint8Array;
export interface PutOptions {
ttlSeconds?: number;
nx?: boolean; // Set only if key does not already exist
xx?: boolean; // Set only if key already exists
}
export interface CacheClient {
// Retrieve value by key, returning null if the key is missing or expired
get(key: string): Promise<CacheValue | null>;
// Store a key-value pair with optional time-to-live and condition flags
put(key: string, value: CacheValue, options?: PutOptions): Promise<boolean>;
// Explicitly delete a key, returning true if removed or false if absent
delete(key: string): Promise<boolean>;
// Atomically increment an integer counter by step (default 1)
incr(key: string, step?: number): Promise<number>;
// Atomically decrement an integer counter by step (default 1)
decr(key: string, step?: number): Promise<number>;
// Compare-and-swap: update key to newValue only if current value matches expectedValue
cas(key: string, expectedValue: CacheValue, newValue: CacheValue): Promise<boolean>;
}Data Structure Operations (Redis-Compatible)
Rich collections support complex access patterns including dictionaries, queues, unique sets, and sorted priority rankings:
export interface RedisDataStructures {
// Hash operations: set field value, retrieve field, or fetch all fields
hset(key: string, field: string, value: string): Promise<number>;
hget(key: string, field: string): Promise<string | null>;
hgetall(key: string): Promise<Record<string, string>>;
// List operations: push to head, pop from tail, or slice index range
lpush(key: string, ...values: string[]): Promise<number>;
rpop(key: string): Promise<string | null>;
lrange(key: string, start: number, stop: number): Promise<string[]>;
// Set operations: add members, retrieve all members, or test membership
sadd(key: string, ...members: string[]): Promise<number>;
smembers(key: string): Promise<string[]>;
sismember(key: string, member: string): Promise<boolean>;
// Sorted Set operations: add member with score, range by rank, range by score
zadd(key: string, score: number, member: string): Promise<number>;
zrange(key: string, start: number, stop: number): Promise<string[]>;
zrangebyscore(key: string, min: number, max: number): Promise<string[]>;
}Common Error Responses
Standard error codes returned by cache proxies and client libraries during operational and connectivity failures:
400 Bad Request: invalid input, missing required fields, or malformed JSON payload 401 Unauthorized: missing or invalid authentication token or API key 403 Forbidden: authenticated caller lacks required permissions for this resource 404 Not Found: requested resource ID does not exist 409 Conflict: duplicate write or version conflict, retry with a unique idempotency key 422 Unprocessable Entity: syntactically valid request failed semantic business validation 429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header 500 Internal Error: unexpected server failure, retry safely with an idempotency key 503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff
Data Model
In-Memory Entry Structure
In-memory cache entries store raw values alongside eviction metadata, timestamp tracking, and doubly linked list pointers:
struct CacheEntry {
char *key; // pointer to key string
char *value; // pointer to value (or embedded for small values)
uint32_t key_length; // 4 bytes
uint32_t value_length; // 4 bytes
uint64_t expiry; // 8 bytes (Unix timestamp, 0 = no expiry)
uint64_t last_access; // 8 bytes (for LRU tracking)
uint8_t lfu_counter; // 1 byte (logarithmic frequency counter)
struct CacheEntry *prev; // 8 bytes (LRU linked list previous pointer)
struct CacheEntry *next; // 8 bytes (LRU linked list next pointer)
// Logical metadata footprint: ~50 bytes before alignment and allocator overhead
};Slab Allocator (Memcached)
Memcached organizes memory into pre-allocated slab pages divided into uniform chunk sizes to eliminate external heap fragmentation:
Slab Classes: Class 1: 64-byte chunks (allocated for values <= 64 bytes) Class 2: 128-byte chunks (allocated for values <= 128 bytes) Class 3: 256-byte chunks (allocated for values <= 256 bytes) Class 4: 512-byte chunks (allocated for values <= 512 bytes) Class 5: 1 KB chunks (allocated for values <= 1024 bytes) ... Class 20: 1 MB chunks (maximum single-item cache payload) Each slab page = 1 MB, divided into uniform chunks of the corresponding class size
Cluster Slot Mapping
Redis Cluster assigns 16,384 slots across master instances with asynchronous replication to dedicated secondary nodes:
Slot Range Master Node Replica Node 0 - 5460 Node A Node A' 5461 - 10922 Node B Node B' 10923 - 16383 Node C Node C'
Fault Tolerance
Resilience mechanisms protect the underlying database from cache stampedes, isolate hot keys, and automate replica promotion during host outages:
| Concern | Solution |
|---|---|
| Node failure | Cluster failure detection detects unresponsive nodes via gossip heartbeats and automatically promotes an eligible replica to primary, while clients refresh their slot routing tables. |
| Data loss on failure | Acceptable because the cache acts as a secondary store, allowing application clients to backfill missing entries directly from the primary database under cache-aside semantics. |
| Cache stampede | Singleflight request coalescing ensures only one worker thread fetches missing data from the database while concurrent readers wait on that in-flight response. |
| Hot key saturation | An in-process L1 cache absorbs high-frequency reads directly in application memory, and hot keys are replicated across multiple shards under partitioned sub-keys. |
| Network partition | Failures remain isolated according to configured quorum policies. Primary shards retaining majority quorum continue serving traffic, while isolated minority shards reject writes to prevent split-brain inconsistencies. |
| Cold start thundering herd | Controlled pre-warming routines populate high-frequency working set keys from query logs before routing live production traffic to newly provisioned cache instances. |
Cache Invalidation Strategies
Selecting the right invalidation mechanism balances data freshness against write latency and recovery complexity:
| Strategy | Mechanism | When to Use |
|---|---|---|
| TTL based | Keys expire automatically after a configured duration | Most common baseline strategy that balances cache freshness with operational simplicity |
| Write Through | Synchronously writes to cache and persistent database in the same operation | Fresh cache state is critical and the application can absorb higher write latency |
| Write Behind | Updates cache immediately and flushes modifications to the database asynchronously | Write-heavy workloads where eventual durability is acceptable, such as view counters |
| Cache Aside (⭐) | Application queries cache first, and on a miss fetches from database and backfills cache | Standard default strategy for general-purpose application workloads |
| Pub/Sub Invalidation | An application change stream, outbox, or CDC pipeline broadcasts invalidations across cluster nodes | Multi-datacenter deployments requiring near-real-time synchronization across regions |
Additional Considerations
Advanced production considerations cover cache warming, operational monitoring, multi-level hierarchy, and resilience against systemic failure modes.
Cache Warming Strategies
Pre-populating cache instances before routing live user traffic prevents cold-start latency surges and avoids immediate database stampedes.
- Deployment Warming: Pre-populate the cache with the top 1,000 to 10,000 most accessed keys identified from access query logs prior to routing user traffic.
- Lazy Warming: Allow the cache to populate organically from database misses, accepting a temporary cold-start latency elevation.
- Hybrid Warming: Proactively warm business-critical catalog and configuration keys while allowing the long tail of low-frequency items to populate lazily.
Monitoring and Telemetry
Comprehensive operational metrics ensure early detection of working set growth, allocator fragmentation, and replication lag:
- Hit Ratio: Target > 95%, computed as
hits / (hits + misses). Sustained drops below 90% indicate a working set shift, premature eviction, or cache key fragmentation. - Memory Utilization: Track memory consumption per node and per slab class to spot fragmentation early.
- Eviction Rate: A sustained increase in eviction rate indicates that the working set is approaching or exceeding provisioned cache capacity or that TTL policies need adjustment.
- Latency Percentiles: Monitor p50, p95, and p99 latencies independently for GET, PUT, and batch pipelined commands.
- Connection Saturation: Track active client TCP socket counts per node to ensure connection pools do not exhaust file descriptors.
- Replication Lag: Measure the time delta between master write completion and replica offset acknowledgment.
Redis vs. Memcached Architecture
Contrasting Redis and Memcached helps engineering teams balance rich collection data structures against lightweight multi-threaded concurrency:
| Feature | Redis | Memcached |
|---|---|---|
| Data structures | Hash, List, Set, Sorted Set, Stream, Bitmap | Only raw strings and blobs |
| Persistence | RDB snapshots and AOF append-only logs | None (pure in-memory storage) |
| Replication | Built-in asynchronous primary-replica replication | None (managed by client-side routing) |
| Clustering | Redis Cluster with 16,384 deterministic hash slots | Client-side consistent hashing rings |
| Threading | Single-threaded core engine with multi-threaded I/O in 6.0+ | Multi-threaded with per-slab locking |
| Memory efficiency | Higher metadata overhead per key (~50 bytes) | Lower overhead using slab class pre-allocation |
| Primary use case | Feature-rich caching, pub/sub, atomic counters, rankings | High-throughput, simple key-value acceleration |
Multi-Level Caching Tiers
Layered caching tiers balance latency and capacity across process memory, distributed nodes, and persistent storage:
Tier 1 (In-Process Cache: Caffeine/Guava): < 0.1 ms latency, ~100 MB heap allocation Tier 2 (Distributed Cache: Redis Cluster): ~1-5 ms latency, ~10-100 GB RAM allocation Tier 3 (Edge Delivery: Regional CDN): ~10-50 ms latency Tier 4 (Primary Persistent Store: Database): ~10-100 ms latency
Cache Penetration, Breakdown, and Avalanche
Recognizing the classic cache failure modes prevents severe downstream database overload during peak load:
| Problem | Description | Solution |
|---|---|---|
| Penetration | Queries target nonexistent keys repeatedly, bypassing cache and striking the database directly | Deploy Bloom filters upstream to intercept missing keys, and cache null values with a short TTL |
| Breakdown | A high-traffic hot key expires, causing thousands of concurrent requests to hit the database | Use singleflight request coalescing with distributed locks, sliding TTLs, or non-expiring keys |
| Avalanche | Large volumes of cached keys share identical TTLs and expire concurrently, triggering massive database surges | Apply randomized TTL jitter (such as ±10%) and stagger cache warming routines across cluster nodes |
Cluster Sizing Math
Worked example for a smaller custom consistent hashing deployment: a 500 GB working set provisioned across nodes with 16 GB of RAM requires 32 nodes at 100% memory utilization. Targeting 70% usable memory to preserve 30% headroom scales the requirement to approximately 46 primary shards. Deploying two replicas per primary shard yields 46 x 3 = 138 Redis instances (replication factor 3). A typical configuration might provision 150 virtual nodes per physical instance. The observed distribution variance depends on hash quality and vnode placement, so the ±5% figure should be treated as an illustrative sizing target rather than a universal guarantee. A single shard failure remaps roughly 1/46 of the key space before the replacement topology stabilizes. Redis Cluster instead rebalances its 16,384 hash slots rather than using this virtual-node calculation.
Related Problems and Core Concepts
In-memory storage and partition routing connect directly to Key-Value Store and API Rate Limiter. Distributed coordination primitives and lock acquisition mechanisms are explored in Distributed Lock Manager. Handling extreme hot key concurrency during traffic surges is covered in Flash Sale System and Like Counter for High-Profile Posts. To master the architectural foundations, review Caching Patterns and Invalidation, Redis Patterns for Interview Systems, Consistent Hashing, Replication, Failover, and Leader Election, Bloom Filter, and System Design Interview Patterns.
Interview Walkthrough
- 25-minute cut
Skip arch50 and arch75 depth unless interviewing for a staff-level role.
- Functional and non-functional requirements and cache-aside baseline (3 min)
- Cache-aside vs write-through trade offs (5 min)
- Consistent hashing for shard placement (6 min)
- LRU eviction and TTL expiration policies (5 min)
- Multi-level cache hierarchy with L1 in-process plus Redis (6 min)
- Default to cache-aside (lazy loading) as detailed in Caching Patterns and Invalidation, because it represents the most battle-tested pattern and simplifies failure recovery.
- Walk through the failure trinity of cache penetration, cache breakdown, and cache avalanche, pairing each failure mode with its concrete mitigation such as Bloom filters, singleflight request coalescing, and randomized TTL jitter.
- Explain shard placement using Consistent Hashing with virtual nodes so adding or removing cluster instances minimizes key remapping to approximately 1/N of total entries.
- Select an eviction policy such as LRU for general workloads or W-TinyLFU for skewed power-law distributions, state a clear target hit ratio above 95%, and monitor eviction rates to detect when working sets outgrow memory.
- Design a layered multi-level cache hierarchy progressing from an in-process L1 cache (sub-0.1ms) to a distributed L2 Redis cluster (1ms) and an edge CDN (10 to 50ms) before hitting the primary database (10 to 100ms), citing latency expectations at each tier.
- Compare Redis and Memcached across data structures, thread concurrency models, memory efficiency, and native replication, guiding your selection based on concrete engineering requirements rather than hype.
- Discuss cache warming on deploy to avoid cold-start latency spikes that trigger avalanches.
- Address the common pitfall of assigning identical TTLs to millions of keys, which causes synchronized expirations that stampede the database.
Engineering Trade-offs
Interview discussions frequently evaluate write strategies, concurrency threading models, sharding schemes, eviction dynamics, cache warming, and thundering herd mitigations. Defend each architectural choice based on concrete operational characteristics.
Write-Through vs Write-Behind vs Cache-Aside: The Critical Pattern Choice
Selecting how writes interact with cache and persistent storage represents the central architectural choice in cache system design:
Cache-Aside (Lazy Loading) ⭐ (Most Common): Read: App checks cache, and on a miss reads DB, writes to cache, and returns value Write: App writes to DB, then invalidates the cache entry or relies on TTL expiration ✓ Simple to implement and maintain ✓ Only requested data is cached (no unused records consuming memory) ✓ Cache failure doesn't prevent reads, although fallback reads are slower ✗ First request is always a cache miss (cold-start penalty) ✗ Data can become temporarily stale if DB updates occur without prompt invalidation Best for: General purpose caching across the majority of web applications Write-Through: Write: App writes to cache, which synchronously writes to DB before returning Read: Served directly from cache when it reflects the latest committed state ✓ Cache reflects the committed write before the operation returns when the cache is the write path ✓ Read operations are consistently fast ✗ Every write incurs extra latency from cache and DB operations ✗ Newly written data may never be read, wasting cache capacity Best for: Read-heavy workloads that need a stronger consistency model and can tolerate higher write latency Write-Behind (Write-Back): Write: App writes to cache, returns immediately, and cache flushes asynchronously to DB ✓ Lowest write latency (acknowledges as soon as memory write succeeds) ✓ Can batch DB writes (1000 cache writes aggregated into 1 bulk DB insert) ✗ DATA LOSS risk: if cache crashes before async flush, uncommitted updates are lost ✗ Complex failure recovery and dual-write reconciliation Best for: Write-heavy workloads where minor data loss is acceptable (activity counters, view metrics) Read-Through: Read: App reads from cache, which auto-fetches from DB on a miss and returns value (Similar to cache-aside, but the cache layer encapsulates the DB fetch logic) ✓ Application client code is simplified ✗ Cache layer must maintain DB connection drivers and query logic (tight coupling)
Why Redis Performs Well With a Mostly Single-Threaded Command Path
A common intuition is that multi-threaded systems are always faster. For in-memory key-value stores, a single-threaded command path eliminates lock synchronization and thread context switching, which can otherwise exceed the 100-nanosecond cost of DRAM lookups:
Intuition: "More threads must always be faster!"
Reality: For many in memory key value workloads, a mostly single threaded data path can be very effective.
Why:
1. No lock contention: Multi-threaded hash maps require synchronization (mutexes, spinlocks, CAS).
Lock contention at 1M operations/sec creates massive CPU thrashing.
Redis avoids mutex contention in its core single threaded command path.
2. In-memory access is ultra-fast: A GET/SET in DRAM takes ~100 nanoseconds.
Thread context switching costs ~1,000 nanoseconds.
At in memory speeds, thread scheduling overhead can exceed the actual data lookup cost.
3. Network I/O is the bottleneck, not CPU core execution:
Parsing a network packet takes ~1 microsecond.
The in-memory dictionary lookup takes ~100 nanoseconds.
In this illustrative network bound workload, the CPU core can remain idle for over 85% of the time while waiting on socket buffers.
Redis 6.0+ hybrid approach: Retains mostly single-threaded command execution,
but offloads network socket I/O to background I/O worker threads.
This unlocks multi-core network throughput without introducing locking overhead.
When multi-threading helps (Memcached advantage):
- Very large values (> 1 KB): CPU time for serialization and payload copying becomes significant
- Very high concurrent connection counts (> 100K): Dedicated I/O worker pools handle socket scaling
- Pure string GET/SET operations: No complex data structures, keeping locking overhead minimalConsistent Hashing: Why It Is Non-Negotiable
Modulo hashing causes catastrophic key redistribution when cluster topology changes. Consistent hashing bounds key movement to roughly 1/N of the keyspace, preventing database collapse during cluster scaling:
Without consistent hashing (simple modulo hashing):
node = hash(key) % N
N = 3 servers: key "user:123" with hash 7 maps to 7 % 3 = 1 (Server 1)
Add a 4th server (N = 4):
key "user:123" with hash 7 maps to 7 % 4 = 3 (Server 3, remapped)
Result: ~75% of all keys are remapped, triggering a massive cache miss storm
All keys miss cache and hit the database simultaneously, risking total database collapse.
With consistent hashing:
Adding Server 4 only moves ~1/N = 25% of keys
The remaining 75% remain assigned to their existing servers without cache misses.
Cost of getting this wrong:
1M keys x 75% remapped x $0.001 per DB query = $750 in emergency DB load
At 1B keys, uncoordinated remapping causes catastrophic cascade failure.Eviction Policy: LRU vs LFU Trade-offs
Selecting between LRU and LFU requires analyzing access distribution skew. LRU preserves temporal recency but is vulnerable to scan pollution, whereas LFU protects persistent working sets across power-law distributions:
LRU (Least Recently Used):
Evicts the key that has gone unaccessed for the longest duration.
Core Assumption: "Recently accessed keys will be accessed again soon."
✓ Excellent for temporal locality (user sessions, recent feed timelines)
✗ Vulnerable to cache pollution: a one-time scan of 1M keys flushes hot items
✗ Ignores frequency: a key accessed 1000 times/hour can be evicted by a one-off request
LFU (Least Frequently Used):
Evicts the key with the lowest cumulative access frequency.
Core Assumption: "Frequently accessed keys represent core working set value."
✓ Superior for power-law skewed workloads (80/20 rule: top 20% of keys serve 80% traffic)
✓ One-time bulk scans do not displace established hot keys
✗ Frequency aging problem: an item viral last week remains cached indefinitely
✗ Cold start penalty: newly written keys start with low frequency and risk instant eviction
Redis LFU implementation:
- 8-bit logarithmic counter with configurable decay so older access history gradually loses influence
- Newly admitted keys start with an initial count = 5 to survive initial admission trials
- Balances frequency protection with active recency decay
Recommendation:
- General purpose workloads: LRU (simple, reliable, minimal CPU overhead)
- Highly skewed access patterns: LFU (some workloads see 5% to 10% higher hit ratios)
- Redis default maxmemory policy: noeviction. allkeys-lru is a common configuration for cache workloadsCache Warming: Cold Start Mitigations
Deploying a cold cache cluster with empty memory exposes the underlying database to an instantaneous 100% miss rate. Structured pre-warming routines protect backend datastores by pre-fetching high-frequency keys before cutover:
Scenario: Deploying a cold cache cluster with empty memory results in a 100% initial miss rate, flooding the database and risking cascading collapse. Solution 1: Lazy warming (passive fallback) Let natural application misses populate cache over time ✗ 5-30 minutes of severe database load and elevated query latency ✗ User-visible latency degradation during cold-start window Solution 2: Pre-warming ⭐ (Recommended) Before routing production traffic to the new cache cluster: 1. Extract top 10,000 most-accessed keys from database query access logs 2. Batch-fetch those records from the primary database and populate the cache 3. Route live traffic to the warmed cache (achieves 70%+ hit rate immediately) Solution 3: Cache replication Instead of cold provisioning, attach new cluster as replicas of active cluster Redis: issue REPLICAOF, perform initial synchronization, and promote to independent master ✓ Zero cold start latency penalty ✗ Requires active source cache cluster to be reachable and sized for sync bandwidth Solution 4: Gradual traffic shift Route 10% traffic to new cache, verify stability, then increase traffic to 25%, 50%, and 100% ✓ Controlled database load during cache population ✗ Requires canary routing infrastructure and prolonged deployment duration
The Thundering Herd and Cache Stampede Mitigations
When a popular hot key expires, thousands of concurrent requests miss the cache simultaneously and overwhelm the database. Request coalescing and probabilistic early expiration collapse redundant queries into a single database fetch:
Timeline:
T=0: Hot key "product:123" cached (TTL = 60s)
T=0-59s: 1,000 requests/sec served directly from cache ✓
T=60s: Key expires
T=60.001s: 1,000 requests arrive simultaneously, all miss cache, and all hit DB concurrently
T=60.050s: Database CPU spikes to 100%, response time jumps to 5 seconds
T=60.100s: Database connection pool exhausts, throwing connection errors to users
T=65s: First DB query completes, cache repopulates, and normal operation resumes
The damage: 5 seconds of total service outage caused by a single key expiring.
Mitigations:
1. Request Coalescing / Singleflight ⭐
The first request becomes the leader for that key, fetches from DB, and populates cache
Remaining requests in the same application process wait on that in-flight work
Only 1 database query is executed instead of 1000 within that process.
For coordination across application instances, add a distributed mutex:
mutex = redis.SET("lock:product:123", "1", NX, EX, 5)
if mutex acquired:
data = db.fetch("product:123")
redis.SET("product:123", data, EX, 60)
return data
else:
deadline = now() + 2s
while now() < deadline:
data = redis.GET("product:123")
if data exists:
return data
sleep(50ms)
return fallback_or_error() // leader may have failed before repopulating cache
2. Stale-While-Revalidate
Cache serves slightly stale data immediately while triggering an asynchronous background fetch
Users receive sub millisecond responses while cache refreshes out of band
3. Randomized TTL Jitter
Instead of uniform TTL = 60s, apply TTL = 60s ± random(0, 10)s
Keys expire asynchronously across time, preventing synchronized herd spikes
4. Never-Expiring Hot Keys with Periodic Background Refresh
Hot keys are provisioned without TTL, while asynchronous cron workers refresh values every 30 seconds
✓ Zero cache misses on hot keys
✗ Consumes constant background polling bandwidthReview
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.