Design a Distributed Cache

Problem Design a distributed cache that scales across many nodes while staying fast, reasonably consistent, and resilient to node loss.

Functional requirements

  • get(key) / set(key, value, ttl) / delete(key).
  • Partition the key space across cache nodes.
  • Evict entries when a node hits its capacity limit.
  • Survive node failure without losing the whole cache.
  • Let clients find the node that owns a key.

Non-functional requirements

  • ~1M ops/sec across ~50 nodes; ~1 TB total cached working set (~20 GB/node), average value ~10 KB.
  • get p99 < 1 ms in-datacenter — the cache exists to be an order of magnitude faster than the origin, or it has no reason to exist.
  • Target hit rate > 90%; every point of hit-rate loss shows up directly as origin DB load.
  • Adding/removing one node should remap ~1/50 (~2%) of keys, not the whole space.
  • Cache loss must be survivable: it is not the source of truth, but a cold start must not stampede the origin into collapse.

Key components

  • Partitioning: consistent hashing over a ring with ~150 virtual nodes per physical node to smooth distribution.
  • Routing: client-side routing (each client holds the ring — one hop, no proxy tier, but every client must agree on membership) or a coordinator/proxy tier (simpler clients, extra hop, another thing to scale).
  • Per-node storage: an in-memory hash map plus an eviction structure — LRU via an intrusive doubly-linked list for O(1) promote/evict, or LFU/TinyLFU when scan resistance matters.
  • Replication: each key on N nodes (primary + N-1 replicas placed at distinct ring positions).
  • Membership/failure detection: gossip or a coordination service, feeding ring updates.
  • Write policy: write-through, write-back, or write-around relative to the origin store.

Deep dives / trade-offs

  • Eviction policy: LRU is simple and matches most access patterns, but one large scan (an analytics job) evicts the entire working set. LFU resists scans but adapts poorly to shifting popularity and needs aging. TinyLFU/W-TinyLFU is the modern default — mention that the policy choice is really about the workload's tail.
  • Write policy: write-through keeps cache and DB consistent at the cost of write latency on every write, and pollutes the cache with data nobody reads. Write-back is fast and coalesces repeated writes but risks data loss on node failure — unacceptable if the cache holds the only copy of a write. Write-around avoids pollution but guarantees a miss on the first read of fresh data.
  • Invalidation: the hard problem. TTL-only is simple but serves stale data for the TTL duration; explicit invalidation is precise but needs every writer to know every cache key affected, and a missed invalidation is a bug that surfaces days later. Discuss versioned keys as an alternative to deletion.
  • Hot keys: consistent hashing balances the key space, not the traffic. One celebrity key saturates its owner node's NIC regardless of ring balance. Mitigations: replicate hot keys across many nodes, add a small client-side local cache in front (with its own coherence problem), or key-split with a random suffix.
  • Consistency: replicas make reads scale and survive failure, but now a set must reach N nodes. Synchronous replication costs write latency; async means a read from a replica can return a stale or deleted value. For a cache, eventual is normally the right call — but say so explicitly, and note when it isn't (a cache holding session or auth state).
  • Thundering herd: on a miss for a hot key, thousands of concurrent requests all hit the origin. Request coalescing / single-flight per key, plus TTL jitter to avoid synchronized expiry.
asked …
LeaderboardSalaryAccount