When an interviewer asks you to design a distributed cache, they are looking for how you balance speed, reliability, and scalability. The conversation usually starts with a vague prompt like “design a cache that can serve millions of requests per second.” Your job is to turn that into a concrete, defensible architecture.

1. Clarify the scope first

Before you draw anything, ask the right questions. Typical clarifying points include:

  • What data is being cached? (e.g., user profiles, product listings, session tokens)
  • Read‑write ratio? Is it mostly reads with occasional writes, or more balanced?
  • Latency expectations? Do we need sub‑millisecond responses, or is a few milliseconds acceptable?
  • Durability requirements? Must the cache survive node failures, or can we tolerate occasional loss?
  • Consistency model? Strong, eventual, or something in between?
  • Geographic distribution? Single data‑center vs. multi‑region deployment.

These answers shape every later decision, so spend a minute here. If the interviewer is vague, propose a reasonable default (e.g., read‑heavy, low latency, eventual consistency) and note that you can adapt later.

2. Define the core API

A clean, language‑agnostic API helps the interviewer see that you understand the contract between clients and the cache.

MethodParametersReturnTypical Use
GET(key)key: stringvalue or nullRetrieve cached data.
SET(key, value, ttl?)key: string, value: any, optional ttlOKInsert or update an entry, optionally with expiration.
DELETE(key)key: stringOKRemove an entry.
INCR(key, delta?)key: string, optional deltaNew numeric valueAtomic counter updates.
GET_MULTI(keys[])keys: list<string>Map of key → valueBatch reads.

Mention that the API can be exposed via gRPC, HTTP/JSON, or a custom binary protocol depending on latency needs.

3. High‑level design diagram (text)

Client --> Load Balancer --> Router --> Shard (Cache Node) <---> Replication Group
                                 |                               |
                                 |                               +--> Persistent Store (optional)
                                 +---> Metrics & Monitoring
  • Load Balancer distributes incoming traffic across routers.
  • Router performs key hashing (e.g., consistent hashing) to locate the responsible shard.
  • Shard is a cache node that stores a partition of the keyspace.
  • Replication Group provides redundancy; each shard has one or more replicas.
  • Persistent Store (optional) backs up writes for durability.
  • Metrics collect latency, hit‑rate, and error stats.

Explain each component briefly as you walk through the diagram.

4. Deep dive: Partitioning and routing

Consistent hashing is the go‑to technique for a distributed cache because it minimizes data movement when nodes join or leave. Outline the steps:

  1. Hash the key to a 64‑bit ring.
  2. Map each node to one or more points on the ring (virtual nodes) to smooth load.
  3. Locate the first node clockwise from the key’s hash – that node owns the key.
  4. On node failure, the next clockwise node takes over.

Mention that you can use a library like MurmurHash3 for the hash function, and that virtual nodes help avoid hot spots.

5. Hard part #1 – Consistency models

Distributed caches rarely give strict linearizability because it hurts latency. Discuss trade‑offs:

  • Strong consistency (read‑after‑write guarantee) requires a quorum protocol (e.g., Raft) or a read‑repair path. Latency rises due to round‑trips.
  • Eventual consistency lets writes propagate asynchronously. Simpler, faster, but stale reads are possible.
  • Read‑your‑writes can be achieved by routing a client’s reads to the same replica that performed the write, using sticky sessions or a client‑side cache token.

A small table can help:

ConsistencyLatency impactWrite complexityTypical use
StrongHigher (extra round‑trip)Coordination neededFinancial transactions
EventualLow (single‑hop)Simple replicationProduct catalog, session data
Read‑your‑writesModerateSlight routing logicUser preferences

Explain why most interviewers expect eventual consistency for a generic cache.

6. Hard part #2 – Eviction and TTL handling

Caches must free space. Two common policies:

  • LRU (Least Recently Used) – easy to implement with a doubly‑linked list; good for hot‑spot workloads.
  • LFU (Least Frequently Used) – better when access patterns are skewed but more memory‑intensive.

TTL (time‑to‑live) can be enforced lazily (on access) or eagerly (background sweeper). Discuss the trade‑off: lazy eviction saves CPU but may keep expired entries longer; eager sweeps keep memory tighter but add periodic work.

7. Failure handling and replication

Explain a typical replication strategy:

  • Primary‑backup: each shard has a leader that accepts writes; followers replicate asynchronously.
  • Quorum reads/writes: a read succeeds after contacting a majority of replicas, guaranteeing a bounded staleness.

On node loss, the router re‑hashes keys to the next live node. Show that the system can continue serving reads from replicas while a new node is provisioned.

8. Observability and operational concerns

Interviewers love to see you think about production:

  • Metrics: hit‑rate, miss‑rate, latency percentiles, eviction count.
  • Logging: request IDs, key hashes, error codes.
  • Health checks: expose /ready and /live endpoints.
  • Capacity planning: monitor memory usage and trigger auto‑scaling.

Mention that a simple dashboard (e.g., Prometheus + Grafana) can surface these metrics.

9. Common follow‑up questions

QuestionTypical angle
"How would you add persistence?"Discuss write‑ahead log, snapshotting, or backing store like DynamoDB.
"What if the cache needs to support transactions?"Explain two‑phase commit or multi‑key locking, and why it hurts performance.
"How do you handle hot keys?"Use request coalescing, sharding by sub‑key, or a separate hot‑key tier.
"Can you make the cache geo‑distributed?"Talk about region‑level routing, read‑only replicas, and WAN latency mitigation.
"What about security?"TLS for transport, token‑based authentication, and ACLs per key prefix.

Prepare concise, high‑level answers; you don’t need to dive into code unless asked.

10. Sample answer snippet (45‑90 seconds)

"Our cache stores key‑value pairs with optional TTL. Clients call GET and SET over a lightweight binary protocol. We partition the keyspace using consistent hashing with virtual nodes to balance load. Each partition has a primary replica that writes synchronously and one or two followers that replicate asynchronously, giving us eventual consistency. For eviction we use an LRU list combined with lazy TTL checks. The router re‑hashes keys on node failure, so the system stays available while a new node joins. Monitoring includes hit‑rate, latency percentiles, and health endpoints, all exported to Prometheus."

11. How to practice this

How to practice this

  1. Sketch the diagram on paper: start with the high‑level components, then add details like hashing and replication.
  2. Explain the API aloud: use Call Assistant to record yourself and get feedback on pacing and clarity.
  3. Run a mock interview: have a friend ask follow‑up questions (e.g., about hot keys) and practice staying on topic while iterating on your design.

FAQ

  1. What is the difference between a cache and a key‑value store? A cache is a short‑lived layer optimized for fast reads, often with eviction policies and TTLs. A key‑value store provides durable storage and may support richer queries.

  2. Why not use a single large cache node instead of sharding? Single nodes become bottlenecks for both throughput and memory. Sharding spreads load, improves fault tolerance, and lets you scale horizontally.

  3. When would you choose strong consistency for a cache? When stale data could cause critical errors—such as financial balances or permission checks—strong consistency outweighs the latency penalty.

  4. How do you prevent a hot key from overwhelming a single shard? Techniques include request coalescing (multiple reads share one backend fetch), splitting the hot key into sub‑keys, or routing hot traffic to a dedicated high‑capacity tier.

Frequently asked questions

What is the difference between a cache and a key‑value store?

A cache is a short‑lived layer optimized for fast reads, often with eviction policies and TTLs. A key‑value store provides durable storage and may support richer queries.

Why not use a single large cache node instead of sharding?

Single nodes become bottlenecks for both throughput and memory. Sharding spreads load, improves fault tolerance, and lets you scale horizontally.

When would you choose strong consistency for a cache?

When stale data could cause critical errors—such as financial balances or permission checks—strong consistency outweighs the latency penalty.

How do you prevent a hot key from overwhelming a single shard?

Use request coalescing, split the hot key into sub‑keys, or route hot traffic to a dedicated high‑capacity tier.

#system design#distributed cache#architecture#consistency#interview prep#a distributed cache