Security Notes
System Design

Design a Distributed Cache

5 min read 7 sections

DifficultyHard | HelloInterview: problem breakdown


Problem Statement

Design a distributed in-memory cache (like Redis Cluster or Memcached) that stores key-value pairs, supports high throughput reads/writes with low latency, and scales horizontally.

๐Ÿ“–

Real-world: The unsung hero of this design is consistent hashing, which came out of a 1997 MIT paper and was originally built to solve exactly this โ€” distributing keys across caching nodes so that adding/removing a node only moves ~1/N of the keys instead of reshuffling everything (the same idea later powering Amazon's Dynamo and DynamoDB). The cautionary security tale for caches: for years tens of thousands of Redis and Memcached instances sat exposed to the internet with no authentication (Redis historically had no auth by default and trusted the network), leading to mass data theft and, worse, Memcached becoming the vehicle for record-breaking DDoS amplification in 2018 โ€” a ~1.3 Tbps attack on GitHub โ€” because an unauthenticated UDP cache reflects huge responses to spoofed victims. The lessons: bind to localhost/private networks, require auth, disable UDP, and never treat "it's just a cache" as "it doesn't need security." Also know the failure modes this design must handle: cache stampede, hot keys, and thundering herd.


Requirements

Functional

  • get(key) โ†’ value
  • put(key, value, ttl)
  • delete(key)
  • Configurable eviction policies (LRU, LFU, TTL)
  • Cluster mode: data distributed across nodes

Non-Functional

  • P99 latency: < 1ms for get/put
  • Throughput: 1M ops/s per cluster
  • High availability: no single point of failure
  • Horizontal scaling: add nodes without downtime
  • Consistency: eventual is acceptable (caches can be stale)

Core Design

Single Node Cache

Start simple: an in-memory hash map with LRU eviction.

HashMap<String, Node> map     โ†’ O(1) lookup
DoublyLinkedList list          โ†’ O(1) LRU move-to-front / evict-from-tail

get(key):
    node = map.get(key)
    if expired: evict, return null
    move node to head of list
    return node.value

put(key, value, ttl):
    if key exists: update value, move to head
    else:
        if at capacity: evict tail
        create node, add to head, add to map

Consistent Hashing for Distribution

Naive approach: hash(key) % N โ€” when N changes (node added/removed), almost all keys need to be remapped. Bad for a live system.

Consistent hashingplace nodes and keys on a conceptual ring (0 โ†’ 2^32). A key belongs to the first node clockwise from its hash position.

Ring: โ”€โ”€โ”€โ”€โ”€ Node A (90ยฐ) โ”€โ”€โ”€โ”€ Node B (180ยฐ) โ”€โ”€โ”€โ”€ Node C (270ยฐ) โ”€โ”€โ”€โ”€
Key K hashes to 120ยฐ โ†’ assigned to Node B (next clockwise)

Add Node D at 150ยฐ:
  Only keys between 120ยฐ-150ยฐ (previously B's) reassign to D
  All other keys: unchanged

Virtual nodeseach physical node occupies multiple positions on the ring (e.g., 150 virtual nodes per physical node). This ensures uniform key distribution even with heterogeneous hardware.

Replication for Availability

Each key is replicated to the next N nodes clockwise on the ring (N = replication factor, typically 3).

Write

primary node receives write, replicates to N-1 replicas

Read

can read from any replica (eventual consistency) OR require quorum (stronger consistency)

Failure

when a node fails, its virtual nodes are taken over by successors; replicas ensure no data loss


Architecture

Client SDK (consistent hash โ†’ pick node)
           โ†“
    โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”
    โ”‚ Cache Node A  โ”‚ Cache Node B     โ”‚
    โ”‚ [key1, key3]  โ”‚ [key2, key4]     โ”‚
    โ”‚ replicas of B โ”‚ replicas of C    โ”‚
    โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜
    โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”
    โ”‚ Cache Node C  โ”‚ Cache Node D     โ”‚
    โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜

Coordinator / Config Server:
  - Tracks live nodes
  - Distributes ring membership to clients
  - Health checks

Client-side vs Proxy-side Routing

Client-side routing (Redis Cluster, Memcached): client library implements consistent hashing, connects directly to correct node. No proxy hop = lower latency. Client must know cluster topology.

Proxy-side routing (e.g., Twemproxy): all clients connect to a proxy; proxy routes to correct node. Simpler clients; proxy is a potential bottleneck.

Eviction Policies

PolicyDescriptionBest for
LRUEvict least recently accessedGeneral cache; temporal locality
LFUEvict least frequently accessedStable hot data; celebrity keys
TTL-basedEvict expired keys proactivelyTime-bounded data (sessions, tokens)
RandomEvict random keyVery simple; approximate LRU at low cost

Active vs lazy expiration

  • Lazy: check TTL on access, evict if expired
  • Active: background thread periodically scans and evicts expired keys (prevents memory leak from never-accessed expired keys)

Key Design Decisions

DecisionChoiceReason
Key distributionConsistent hashing with virtual nodesMinimize remapping on membership change
Replication3 replicas, asyncAvailability without strong consistency overhead
ConsistencyEventualCaches are acceptable to be slightly stale
EvictionLRU + TTLMost common access pattern
RoutingClient-sideLower latency, no proxy bottleneck
SerializationMessagePack or ProtobufCompact binary; faster than JSON

Security Considerations

ThreatMitigation
Cache poisoningValidate data written to cache; cache should not be the source of truth for auth decisions
Cache key enumerationUse opaque/hashed cache keys, not predictable patterns
Sensitive data in cacheEncrypt sensitive values (PII, tokens); set short TTL
Unauthorized cache accessNetwork-level access control (VPC, security groups); auth tokens for cache access
Cache side-channelCache timing attacks can leak whether a key exists; use consistent-time responses
Denial of serviceRate limit per client on cache access; limit value sizes; reject excessively large values
Cache stampede on invalidationProbabilistic early expiration; mutex lock on miss; staggered TTLs

AuthenticationRedis in production should always have requirepass or ACL-based auth, and should never be exposed on a public interface without authentication.


Interview Tips

Start with the single-node LRU

walk through the data structure before scaling.

Consistent hashing is the key concept

spend time on why % N is bad and how the ring fixes it.

Virtual nodes = uniform distribution

don't skip this, it's why consistent hashing is practical.

Replication factor and quorum reads

connect this to CAP theorem: you're trading consistency for availability and partition tolerance.

Common follow-up

"how do you handle cache invalidation at scale?" โ€” event-driven invalidation via message queue, vs TTL, vs version keys.