Learning pathsA
Question Breakdowns

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

Contract / pseudocode
GET key -> {value,version} or miss
SET key value ttl [expectedVersion]
DELETE key
MGET keys -> per-key results with bounded request size

Data model and access paths

Contract / pseudocode
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

Three architecture decisions for Distributed Cache, including the pressure each introduces.
Scroll to inspect the diagram, or open it at full size.

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

Connected responsibilities for Distributed Cache. Trace the authoritative and derived paths separately.
Scroll to inspect the diagram, or open it at full size.

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

Distributed Cache: baseline request paths.
Scroll to inspect the diagram, or open it at full size.
  1. Client routing table → key owner → in-memory entry.
  2. TTL or eviction → miss → authoritative origin.
  3. Origin result → conditional cache fill.

Evolve the design under load

Distributed Cache: additional scaling and recovery paths.
Scroll to inspect the diagram, or open it at full size.
  1. Virtual nodes → weighted placement → replication.
  2. Membership epochs → controlled key migration.
  3. 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

MechanismBenefitCost
Consistent placementLimited remappingMembership coordination
Replica readsHot-key reliefStaleness and invalidation
Request coalescingProtects origin on missesPer-key in-flight state
Bounded fallbackSurvives cluster lossSome 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.

A ring placement example: adding a token changes one ownership interval, not every key.
Scroll to inspect the diagram, or open it at full size.
One key during ring migration. 7 / Old owner token 60 / Source version accompanies cached value; 8 copying / New token 50 receives copied values / Copy cannot overwrite newer source version; 8 active / Router targets new owner; old rejects stale writes / Retry with current topology; Old replica promoted without latest write / May contain stale cache value / Respect allowed staleness or fetch origin
Scroll to inspect the diagram, or open it at full size.

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.

python
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 fresh

Negative 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.

One origin fill for many concurrent misses. Clients to Cache router: Many requests for viral key K; Cache router to Refill owner: Join one bounded in-flight fill for K; Refill owner to Origin: One source read under global origin budget; Origin to Refill owner: Return source version 18; Refill owner to Cache router: CAS fill if version not older; release waiters; Cache router to Clients: Return value or declared stale/error policy
Scroll to inspect the diagram, or open it at full size.

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.

Technical references

Redis replication guarantees.

PostgreSQL transaction isolation and concurrent updates.

Your study notes