- Published on
- Views
Design a Distributed Cache System Design Interview Guide
- Authors

- Name
- Javed Shaikh
← 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/42is 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:
- Stores key-value data in memory across several machines
- Lets any app server read and write those keys
- Evicts old or unused data
- 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
| Requirement | What it means |
|---|---|
| GET / SET / DELETE | Basic key-value API used by app servers. |
| TTL | Every key can expire. Stale data should disappear. |
| Shared across servers | Two Spring Boot pods see the same value for the same key. |
| Eviction | If memory is full, remove something (usually LRU). |
| Invalidation | When the database row changes, the cache key can be deleted. |
| Large cluster | Data 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
| Requirement | Why it matters |
|---|---|
| Low latency | Cache GET should be sub-millisecond to a few milliseconds. |
| High hit ratio | Misses become database load. Aim high for hot keys. |
| Availability | If one cache node dies, the rest should still serve. |
| Memory efficient | RAM is the expensive part. |
| Predictable eviction | Full cache should not pause the whole cluster. |
| Safe fallback | Cache 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
GET product:42from cache- If hit, return
- If miss,
SELECTfrom DB SET product:42with TTL- Return to client
PUT /api/products/42
- Update DB
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
| Key | Value | TTL | Why |
|---|---|---|---|
product:{id} | product JSON | 5–30 min | Hot catalog reads |
user:{id}:profile | profile JSON | 10 min | Many screens need it |
session:{token} | session blob | session length | Fast auth |
feed:{userId}:page:{n} | list of ids | 30–60 sec | Short-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
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):
- Client calls the app.
- App computes the key.
- Cache GET.
- Hit → return.
- Miss → load DB → SET cache with TTL → return.
Write flow:
- Client updates data.
- App writes the database first.
- App deletes (or updates) the cache key.
- 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:
- TTL only — simple, stale until expiry
- Delete on write —
DELafter DB update - Version in the key —
product:v7:42after a schema change - 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
| Bottleneck | What happens | What you do |
|---|---|---|
| Hot key | One shard melts | Local cache, replicate key, CDN |
| Low hit ratio | DB overload | Cache the right keys. Check TTL. Warm after deploy. |
| Stampede | Many misses on one key | Lock / single-flight. Slightly randomized TTL. |
| Large values | Network and RAM waste | Cache ids, not whole graphs. Compress. |
| Invalidation storms | Bulk edits delete millions of keys | Version prefixes. Cache objects, not huge lists. |
| Unbalanced shards | One node at 90% RAM | Virtual 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.
15. Internal Links
Related JavaThoughts reading:
- System Design Interview Preparation
- Design a URL Shortener
- Design a Rate Limiter
- Caching
- Consistent hashing
- Java
- Spring Boot
- Distributed Systems
- Microservices
- 20 system design concepts
Next in this series: Design a News Feed.
