Javathoughts Logo
Javathoughts
Published on
Views

Design a Distributed Cache System Design Interview Guide

Authors
  • avatar
    Name
    Javed Shaikh
    Twitter

← System Design Interview Preparation

This guide walks through Design a Distributed Cache the way you would in a backend or Java interview. The numbers are interview estimates. They help you show your thinking. They are not a production capacity plan.


1. Problem

A distributed cache is a shared memory store in front of a slower database. Many app servers read and write the same cache cluster.

Example:

  • Product page GET /products/42 is requested thousands of times per second
  • The database row barely changes
  • App servers ask Redis/Memcached first
  • On a miss, they load the database and put the result in cache with a TTL

You use this in a URL shortener redirect path and in almost every read-heavy Java service.

In this design we are building a cache that:

  1. Stores key-value data in memory across several machines
  2. Lets any app server read and write those keys
  3. Evicts old or unused data
  4. Survives the loss of one cache node without a full outage

We are not rebuilding Redis from scratch in 45 minutes. We are designing how the application and the cache cluster work together, plus the cluster ideas interviewers expect: sharding, TTL, invalidation, hot keys.


2. Functional Requirements / FR

RequirementWhat it means
GET / SET / DELETEBasic key-value API used by app servers.
TTLEvery key can expire. Stale data should disappear.
Shared across serversTwo Spring Boot pods see the same value for the same key.
EvictionIf memory is full, remove something (usually LRU).
InvalidationWhen the database row changes, the cache key can be deleted.
Large clusterData is split across nodes so one box is not the whole cache.

Out of scope for a 45-minute interview:

  • Full Redis Cluster protocol details
  • Disk persistence internals
  • A query language
  • Transactional SQL in the cache

Confirm we are designing cache-aside as the main pattern. Mention read-through and write-through briefly.


3. Non-Functional Requirements / NFR

RequirementWhy it matters
Low latencyCache GET should be sub-millisecond to a few milliseconds.
High hit ratioMisses become database load. Aim high for hot keys.
AvailabilityIf one cache node dies, the rest should still serve.
Memory efficientRAM is the expensive part.
Predictable evictionFull cache should not pause the whole cluster.
Safe fallbackCache down should not take the product down forever.

A good interview sentence: the database is the source of truth. The cache is a fast, possibly stale copy.


4. Back-of-the-Envelope Calculation

Say these assumptions out loud.

Traffic assumptions

  • 200 million page/API reads per day
  • 10 million writes / updates per day
  • Read/write ratio ≈ 200M / 10M = 20 : 1

This is a read-heavy system. Cache is the right tool.

QPS

Average read QPS:

200,000,000 / 86,400 ≈ 2,300 QPS

Peak 5× → design for around 12,000 read QPS.

Average write QPS:

10,000,000 / 86,400 ≈ 116 QPS

Peak writes ≈ 600 QPS. Writes matter because they trigger invalidation.

If cache hit ratio is 90%:

Cache GETs ≈ 12,000 / sec
Database reads ≈ 1,200 / sec  (the 10% misses)

Without cache, the database would see all 12,000 QPS. That is why the cache exists.

Storage / memory estimate

Assume we cache 20 million hot objects. Average value size 2 KB (JSON product, user profile, session).

20,000,000 × 2 KB ≈ 40 GB

Add ~30% for Redis overhead, fragmentation, and replicas:

40 GB × 1.3 ≈ 52 GB

A cluster of 3 nodes × 20 GB usable RAM, or 6 nodes × 12 GB, is a reasonable interview picture. Plus one replica per shard if you need HA.

Server estimate

App servers: a Java service that mostly hits cache might do 2,000–4,000 QPS per instance.

Peak 12,000 QPS / 3,000 ≈ 4 app servers

Use 6 for headroom. Cache nodes are sized by memory, not only by QPS. 12,000 GETs/sec is light for Redis. Memory is the real budget.

These are interview estimates, not exact production numbers.


5. APIs

The cache API is small. App servers call it. Users do not.

Get

GET /cache/{key}

{
  "key": "product:42",
  "value": { "id": 42, "name": "Notebook", "price": 499 },
  "ttlSeconds": 218
}

Miss: 404 from the cache client, then the app loads the DB.

Set

PUT /cache/{key}

{
  "value": { "id": 42, "name": "Notebook", "price": 499 },
  "ttlSeconds": 300
}

Delete (invalidation)

DELETE /cache/{key}

Used after a product update.

In Java interviews you can say: Lettuce/Jedis GET, SET key value EX 300, DEL. No custom HTTP is required. Showing the operations is enough.

App-level product API (how cache is used)

GET /api/products/42

  1. GET product:42 from cache
  2. If hit, return
  3. If miss, SELECT from DB
  4. SET product:42 with TTL
  5. Return to client

PUT /api/products/42

  1. Update DB
  2. DEL product:42

That is cache-aside.


6. Data Model

Cache is a key-value store, not a relational schema. Still, write the key design. Bad keys cause bugs.

Key design

KeyValueTTLWhy
product:{id}product JSON5–30 minHot catalog reads
user:{id}:profileprofile JSON10 minMany screens need it
session:{token}session blobsession lengthFast auth
feed:{userId}:page:{n}list of ids30–60 secShort-lived, changes often

Rules:

  • Include a namespace (product:) so keys do not clash
  • Include a version if the JSON shape changes (product:v2:42)
  • Keep keys short
  • Never use a huge unhashed list as a key

Metadata on a cache node

Each node tracks:

  • key → value bytes
  • expiry time
  • last access time (for LRU)
  • approx memory used

Database remains the source of truth

The product table does not go away. Cache can vanish. The app must still know how to fill it.


7. High-Level Design

The app service talks to a cache cluster first. Misses go to the database. The result is stored back in cache.

Distributed Cache architecture

Distributed Cache architectureThe app reads from the cache cluster first. On a miss it loads from the database and writes the value back into cache with a TTL.get / setcache missstore result👤Client⚙️App Service🧠Cache Cluster🗄️Database
The app reads from the cache cluster first. On a miss it loads from the database and writes the value back into cache with a TTL.

Components:

  • Client: browser or mobile app.
  • App Service: Java/Spring Boot business logic. It owns the cache-aside flow.
  • Cache Cluster: several Redis/Memcached nodes. Keys are partitioned with consistent hashing.
  • Database: source of truth.

Read flow (cache-aside):

  1. Client calls the app.
  2. App computes the key.
  3. Cache GET.
  4. Hit → return.
  5. Miss → load DB → SET cache with TTL → return.

Write flow:

  1. Client updates data.
  2. App writes the database first.
  3. App deletes (or updates) the cache key.
  4. Next read refills cache.

Why not write cache first? If the DB write fails, cache would lie. For most product data, DB first, then invalidate is the safe interview default.


8. Deep Dives

Cache-aside pattern

The application is in charge.

  • Read: cache then DB
  • Write: DB then delete cache

Pros: cache stays simple. You choose what to cache.
Cons: first request after invalidation is slow. Two app servers can both miss and both hit the DB (cache stampede). For hot keys, add a single-flight lock or a short negative lock.

This is the pattern you should draw first.

Read-through / write-through (brief)

Read-through: the cache library loads the DB on miss for you. App only calls GET.

Write-through: app writes to the cache, cache writes to the DB.

Useful when you want one library to enforce policy. Less common in Spring interviews than cache-aside. Mention them so the interviewer knows you have heard of them, then go back to cache-aside.

Write-behind: cache writes to DB later. Faster writes, risk of loss if cache dies. Rare for money or inventory.

TTL

TTL is how you accept eventual consistency.

  • Short TTL: safer, more DB load
  • Long TTL: faster, more stale reads

Examples:

  • Stock count: seconds, or do not cache
  • Product title: minutes to hours
  • Country list: a day

Always set a TTL. A key with no TTL plus a missed invalidation lives forever.

Eviction

When RAM is full:

  • LRU: remove least recently used. Default interview answer.
  • LFU: remove least frequently used. Helps if some keys are old but still hot.
  • TTL expiry: still needed even with LRU
  • No eviction (volatile-lru): only keys with TTL are evicted

Say: size the cluster so working set fits. Eviction is the safety net, not the plan.

Cache invalidation

This is the hard part. “There are only two hard things in computer science…” exists for a reason.

Strategies:

  1. TTL only — simple, stale until expiry
  2. Delete on write — DEL after DB update
  3. Version in the key — product:v7:42 after a schema change
  4. Pub/sub — writer publishes “invalidate product:42”, all app local caches drop it

If you have local Caffeine + Redis, you must invalidate both. Redis is shared. Caffeine is per JVM.

Consistent hashing

You have N cache nodes. You must decide which node holds product:42.

Naive hash(key) % N breaks when you add a node: almost every key moves, cache empties, database melts.

Consistent hashing places nodes on a ring. A key maps to the next node on the ring. When you add a node, only a slice of keys move. Virtual nodes (many points per machine) keep load even.

Interview sentence: consistent hashing exists so scaling the cluster does not flush the whole cache.

See also consistent hashing.

Hot keys

One celebrity product, one viral short code, one logged-in-home-page blob.

Problems:

  • One Redis CPU / network card saturates
  • One shard’s memory fills with a huge value

Mitigations:

  • Replicate that key on several nodes and pick randomly
  • Local in-process cache for the hottest keys (1–2 seconds)
  • Split a huge value
  • CDN for public pages

Name hot keys in every cache interview.

Redis vs Memcached style

  • Redis: data structures, TTL, persistence option, Lua, pub/sub. Default for Java shops.
  • Memcached: simple LRU memory, very fast, no fancy types.

If the interviewer says “design Memcached,” focus on hashing, LRU, and multithreaded memory. If they say “Redis,” you can also mention sorted sets for leaderboards, but do not turn this into a Redis feature list.


9. Bottlenecks

BottleneckWhat happensWhat you do
Hot keyOne shard meltsLocal cache, replicate key, CDN
Low hit ratioDB overloadCache the right keys. Check TTL. Warm after deploy.
StampedeMany misses on one keyLock / single-flight. Slightly randomized TTL.
Large valuesNetwork and RAM wasteCache ids, not whole graphs. Compress.
Invalidation stormsBulk edits delete millions of keysVersion prefixes. Cache objects, not huge lists.
Unbalanced shardsOne node at 90% RAMVirtual nodes. Split fat keys.

The first bottleneck to name is hot keys plus cache misses hitting the database together.


10. Tradeoffs

Cache-aside vs write-through

  • Cache-aside: simple, app owns correctness, extra code in every service.
  • Write-through: always warm after writes, extra write latency, cache is in the critical path.

Pick cache-aside for most interviews.

TTL vs explicit invalidation

  • TTL only: you will serve stale data
  • Invalidation only: a missed delete serves stale data forever
  • Both: invalidate on write, TTL as a backstop

Consistency vs speed

A cache by definition can be stale. If the interviewer asks for strong consistency, say: do not cache that field, or use a version check against the DB for checkout.

Memory vs hit ratio

More RAM → more keys → better hits → less DB. Cost goes up. Size for the working set, not the whole database.


11. Failure Modes

Cache node down
Consistent hashing moves keys to neighbors, or you fail over to a replica. Miss ratio jumps. Database must have headroom. This is why you never size the DB for “cache always works”.

Whole cache cluster down
App goes to the database. Autoscale reads if you can. Shed non-critical cacheable traffic. Circuit-break expensive queries.

Stale read after update
Writer updated DB but DEL failed. TTL will fix it. Retry delete. For prices at checkout, read DB.

Thundering herd after deploy
Empty cache. Warm top keys. Use randomized TTL so they do not expire together.

Split brain with local caches
JVM A has old product, Redis has new. Keep local TTL very short (1–5 seconds) or subscribe to invalidation.

Hot key eviction
LRU kicked out the celebrity key because a scan loaded junk. Pin hot keys or use a separate tiny cache for them.


12. Interview Answer in 10 Minutes

Here is a version you can speak out loud.

"I would design a distributed cache as a cache-aside layer in front of the database. The database stays the source of truth. App servers check Redis first.

Assume 200 million reads a day and 10 million writes. That is about 2,300 read QPS, around 12,000 at a 5× peak. If we hit cache 90% of the time, the database only sees about 1,200 QPS at peak. I would cache around 20 million hot objects of 2 KB, which is about 40 GB of data, plus overhead, so a small Redis cluster.

The API is GET, SET with TTL, and DELETE. On read: get from cache, on miss load DB and fill cache. On write: update DB, then delete the key. I always set a TTL so a missed invalidation cannot live forever.

To split keys across nodes I would use consistent hashing with virtual nodes, so adding a machine does not flush the whole cache. For eviction I would use LRU. For hot keys I would add a tiny local cache or replicate that one key.

If Redis is down, we fall back to the database and protect it with timeouts. If one node dies, hashing moves a slice of keys and miss ratio goes up for a while.

That is the design: cache-aside, TTL plus delete on write, consistent hashing, and a plan for hot keys."

Practice this until it is under 10 minutes.


13. Interview Talking Points

  • Database is source of truth. Cache can vanish.
  • Cache-aside as the default pattern.
  • TTL + delete on write.
  • Consistent hashing when the cluster grows.
  • Hot keys and stampede are the real production bugs.
  • Local cache + distributed cache if JVM heap is used.
  • Monitor hit ratio, memory, evictions, and DB QPS.
  • Never cache data you cannot afford to serve stale without a plan.

14. Follow-up Questions

How do you handle a celebrity key?
Local 1-second cache, replicate the key, or put a CDN in front if it is public HTML.

What is cache stampede?
Many servers miss together and all hit the DB. Use a lock for that key, request coalescing, or slightly random TTLs.

How does consistent hashing help?
Only a fraction of keys move when you add or remove a node, so you do not empty the cache.

Redis or Caffeine?
Caffeine is per process, Redis is shared. Use Caffeine for tiny ultra-hot data, Redis for shared state. Invalidate both.

Can I cache inventory?
Only with a very short TTL, or not at all for the final checkout read. Wrong stock is a business bug.

Other useful probes: write-behind risk, compression, and multi-region cache.


Related JavaThoughts reading:


Next in this series: Design a News Feed.