Distributed Cache
Design a cache cluster that serves key/value reads through node failures and membership changes.
Interview scope and guarantees
Provide get, set, delete, TTL and optional compare-and-set over partitioned keys. State whether the cache is disposable or a system of record; this design assumes authoritative data exists elsewhere. Define eviction and miss behavior under node loss.
Capacity worksheet
Assume 100 million entries averaging 500 bytes including logical metadata: 50 GB raw. Two copies require 100 GB before allocator overhead and free-space headroom. At 1 million requests/s, a 1% miss rate still creates 10,000 origin reads/s; losing a shard can abruptly raise that.
Concrete API contract
GET key -> {value,version} or miss
SET key value ttl [expectedVersion]
DELETE key
MGET keys -> per-key results with bounded request sizeData model and access paths
membership(epoch,virtual_token,node_id,state)
entry(key,value,version,expires_at)
replica_assignment(token,primary,secondary)
request_budget(client,origin_fallback_limit)Evolve a solution and explain each change
Figure — Three architecture decisions for Distributed Cache, including the pressure each introduces.
Step 1: Route keys
Use a defined hash topology to find a key owner with TTL and byte limits. Rebalancing and node loss create misses and stale ownership.
Step 2: Add replicas carefully
Replicate according to a stated consistency and failover model. An acknowledged write can be absent after asynchronous promotion.
Step 3: Protect the origin
Use bounded refill, hot-key replication, and versioned topology transitions. Cache availability must not imply authoritative business ownership.
Responsibility overview
Figure — Connected responsibilities for Distributed Cache. Trace the authoritative and derived paths separately.
Worked end-to-end scenario
Client C writes key k to primary P under topology epoch 7. P acknowledges before replica R receives it, then fails. Promotion can return the previous value or a miss. For a derived cache this may be acceptable if the origin is authoritative and stale-read limits are met. During rebalance to epoch 8, clients need a coherent routing update or controlled forwarding; blindly writing both owners can produce conflicting versions. A popular key should be replicated near readers, while a burst of misses needs single-flight or bounded origin requests. Expiring many keys together can overwhelm the database even when cache nodes are healthy.
Why these access paths matter
Keys map to virtual-node ranges or another explicit partition function. Entries carry expiration, size, and optional source version. Choose eviction by workload rather than assuming LRU fits all keys. Store topology separately from values and make clients handle epoch changes. A cache lock without fencing cannot safely protect durable ownership in an external system.
Build the baseline first
- Client routing table → key owner → in-memory entry.
- TTL or eviction → miss → authoritative origin.
- Origin result → conditional cache fill.
Evolve the design under load
- Virtual nodes → weighted placement → replication.
- Membership epochs → controlled key migration.
- Hot-key replicas → local coalescing → origin protection.
Defend the hardest decision
Consistent hashing reduces ownership movement when membership changes, but does not replicate values or solve hot keys. A client needs a versioned routing map and redirection behavior during migration. If the cache is disposable, a migration can tolerate misses and refill; if it promises durable acknowledged writes, it needs a far stronger storage and replication design. State that boundary before discussing algorithms.
Failure and recovery analysis
A primary fails and its replica has stale values. Serving them may violate an application’s freshness contract. Options are invalidate-on-failover, version validation, or accepting bounded staleness explicitly. Eviction should not remove unrelated durable state. A stampede after shard loss requires origin admission control and request coalescing, not merely more cache instances.
Security and privacy boundary
Isolate tenant key namespaces, authenticate clients, cap values and batch sizes, and never assume private network placement is sufficient authorization.
Interview follow-ups with reasoning
Question: How do you split a hot key?
Show answer and explanation
Answer: Replicate reads or redesign its value; hashing the same key again does not divide traffic.
Question: How do you test rebalancing?
Show answer and explanation
Answer: Track moved keys, miss spikes and origin load under membership churn.
Question: What is the cost of LRU?
Show answer and explanation
Answer: Metadata updates and coordination; approximate policies may be preferable.
Operate and verify the design
Hit rate by bytes, eviction count, origin fallback rate and shard-rebalance load.
Remove one cache shard and verify origin admission prevents a cascading database outage.
A second scenario to test transfer
A node owns 10% of the ring but holds 60% of request volume because one key is viral. Moving more ordinary keys will not solve the hot key. Replicate that read, cache it near consumers, or split its representation if semantics allow. Measure request heat and bytes by key, not only key-count balance.
A cache node fails and its misses overwhelm the database. Which changes reduce the immediate blast radius?
Show answer and explanation
Answer: Throttle concurrent origin fills, coalesce requests for identical hot keys, serve stale values only where allowed, and shed optional traffic. Then restore cache capacity and warm critical keys gradually. Adding application retries or unbounded database connections amplifies the outage.
Compare alternatives
| Mechanism | Benefit | Cost |
|---|---|---|
| Consistent placement | Limited remapping | Membership coordination |
| Replica reads | Hot-key relief | Staleness and invalidation |
| Request coalescing | Protects origin on misses | Per-key in-flight state |
| Bounded fallback | Survives cluster loss | Some requests must wait or fail |
A design-changing exercise
A lock key expires while its holder pauses. Can the resumed holder still safely write the database?
Show answer and explanation
Answer: Not based on the expired lock alone. Durable writes need a fencing token or another authoritative conditional ownership check.
Design workshop: placement, storage, and refill protection
Core scope is disposable key/value caching with TTL and bounded origin fallback. Durable queue semantics and authoritative payment state are excluded. Assume p95 hit latency under five milliseconds within a region and a disclosed stale-read bound where replicas are used. A cache write acknowledgment is not an origin commit.
Inside each node, a hash table maps keys to entries containing value, source version, expiry and allocation metadata. A size-aware allocator can use slabs/classes; fragmentation and metadata count against the memory budget. Choose approximate LRU or LFU with bounded sampling rather than an unbounded exact global linked structure. Evict expired entries first when encountered and use background bounded expiry work. GET treats expired entries as absent even if physical deletion has not run. SET-if-version/CAS is atomic at the shard owner.
A ring example uses tokens at 20,60,90 and clockwise successor ownership. Key hash 45 belongs to token 60. Adding token 50 moves only the interval (20,50] to the new owner, including key 45; it does not move every key. Virtual nodes improve placement balance but cannot split one hot key's request traffic automatically.
Figure — One key during ring migration.
Membership authority uses consensus-backed epochs or a single durable topology coordinator. Activate a new epoch only after the old owner is fenced under the chosen lease contract; partitioned clients cannot promote themselves. Since this cache is derived, missing migrated values can refill, but two refill writers must still compare source versions. Replication mode determines the tolerated loss/staleness; asynchronous failover may discard acknowledged cache updates.
value = cache.get(key)
if value and valid_for_request(value): return value
with bounded_singleflight(key):
value = cache.get(key)
if value: return value
if not origin_budget.try_acquire(): return stale_or_503()
fresh = origin.read(key) # fixed timeout and bounded connections
cache.set_if_source_version_not_older(key, fresh, jittered_ttl())
return freshNegative cache entries use shorter TTL and distinguish absent from permission denied and transient origin failures. Never cache an authorization denial under a key shared by every user. Hot values can use near caches/read replicas within their stated freshness contract; a strict revocation path may bypass them.
Figure — One origin fill for many concurrent misses.
For ten million entries averaging 500-byte values, 50-byte keys and 100-byte metadata, primary memory is about 6.5 GB before allocator fragmentation. Two replicas increase storage copies, not necessarily useful independent hot-key capacity. Measure bytes, evictions, hit ratio by key heat, origin refill concurrency and topology epoch errors. Cold cluster recovery warms gradually; unbounded misses defeat the purpose of caching.
Exercise: A node owns 10% of keys but 60% of traffic. Will moving random cold keys fix overload?
Show answer and explanation
Answer: No. Identify the hot key, use allowed replica/near-cache reads or partition its representation if semantics permit, and bound refill traffic. Membership balance and access heat are separate dimensions.
Redis replication documents acknowledgment/replication limits relevant to this cache design.