NEW: ML Mock & Coaching now available

Questions

Distributed Cache

MetaGoogle

Design a scalable distributed cache system supporting key-value operations, TTL expiration, eviction policies (LRU/LFU), consistent hashing, replication, and configurable consistency models.

40 min read

Challenge

Think Beyond the Happy Path

Distributed caches aren't just about fast key-value lookups — they're about handling hot keys, automatic failover, and choosing the right consistency model for your use case.

Before diving into the material, take a moment to ask yourself:

  • Do you know how to achieve sub-millisecond latency at scale — what data structures and concurrency patterns enable O(1) lookups across hundreds of nodes?
  • Do you know how to handle hot keys that receive millions of reads per second — how do you distribute load without creating consistency problems?
  • Do you know how to design replication and automatic failover — what happens when a cache node crashes mid-operation, and how do you avoid data loss?
  • Do you know when to choose strong vs. eventual consistency — and how to implement atomic operations like increment and compare-and-set in a distributed setting?

You don't need to answer all of these right away. A Senior/Staff+ Engineer doesn't stop at basic functionality — they anticipate edge cases, design for resilience, and push for production-grade reliability.

But great systems start with great questions. What would you ask next?

Problem Statement

Designing a distributed cache involves balancing multiple factors—availability, consistency, scalability, latency, and cost—so that you can handle large amounts of data efficiently across multiple nodes.

When to Use a Distributed Cache
  • High Read Volume: You have a workload with many read operations and fewer writes.
  • Data Reuse: Certain data is accessed repeatedly by different services or nodes.
  • Latency Requirements: You need to reduce round trips to a database or external service to improve response time.

When NOT to Use a Distributed Cache

  • No Reuse: Data is seldom reused or is always unique (e.g., real-time streams).
  • Write-Heavy Workloads: When writes dominate, cache invalidation overhead can outweigh benefits.
  • Strong Consistency Required: When you cannot tolerate any staleness in data.

Before diving into the design, consider:

  • How do you ensure sub-millisecond latency at scale?
  • How do you handle hot keys that receive millions of requests per second?
  • What happens when a node fails mid-operation?
  • How do you rebalance data without service disruption?

A Senior/Staff+ engineer designs for these edge cases from day one.

Functional Requirements

  • FR1: Basic Key-Value Operations

    Enable set, get, update, and delete operations, optionally supporting bulk operations for efficiency.

  • FR2: TTL Support

    Support TTL (Time-To-Live) for automatic expiration of key-value pairs.

  • FR3: Configurable Eviction Policies

    Implement eviction policies (LRU, LFU) for managing storage limits when memory is exhausted.

  • FR4: Partitioning & Rebalancing

    Provide partitioning options (e.g., consistent hashing) with automatic data and load rebalancing.

  • FR5: Caching Strategies

    Support write-through, write-behind, and read-through caching patterns.

Info

Bonus: Even though partition and replication are usually considered non-functional requirements, due to the nature of this particular problem, we can make consistency and partition configuration part of the functional requirements.

Non-Functional Requirements

  • NFR1: Performance — Ensure low-latency operations (<10ms for most requests) with efficient concurrency handling.
  • NFR2: Scalability — Support horizontal scalability to add/remove nodes and handle increasing data and traffic.
  • NFR3: Fault Tolerance — Provide replication (master-slave or peer-to-peer) and automatic failover for high availability. Offer durability with persistence and disaster recovery via multi-AZ/region replication.
  • NFR4: Consistency — Allow configurable consistency (strong/eventual) based on use case requirements.
  • NFR5: Atomic Operations — Enable atomic operations like increment and compare-and-set for data integrity.
Requirements Summary
NameDescription
FR1Basic Key-Value Operations | Set, get, update, delete with optional bulk operations.
FR2TTL Support | Automatic expiration of key-value pairs with configurable time-to-live.
FR3Configurable Eviction Policies | LRU, LFU policies for managing memory limits.
FR4Partitioning & Rebalancing | Consistent hashing with automatic data redistribution.
FR5Caching Strategies | Write-through, write-behind, and read-through patterns.
NFR1Performance | Sub-10ms latency for most requests with efficient concurrency.
NFR2Scalability | Horizontal scaling to add/remove nodes seamlessly.
NFR3Fault Tolerance | Replication, automatic failover, multi-AZ/region disaster recovery.
NFR4Consistency | Configurable strong or eventual consistency models.
NFR5Atomic Operations | Increment, compare-and-set for data integrity.
Info

Monitoring is Crucial: For infrastructure-style questions, add monitoring for each functional requirement: hit/miss ratios, latency percentiles, memory usage, eviction rates, and replication lag.

Below-the-Line Requirements

  • Secure access with authentication, authorization, and TLS/SSL encryption.
  • Support backup/restore tools and persistent snapshots for reliability and recovery.

API Design

Based on our functional requirements, we need to expose core operations through an API:

// Basic Operations
SET(key, value, ttl?) -> OK | ERROR
GET(key) -> value | NULL
DELETE(key) -> OK | NOT_FOUND

// Bulk Operations
MSET([(key1, value1), (key2, value2), ...]) -> OK | PARTIAL_ERROR
MGET([key1, key2, ...]) -> [value1, value2, ...]

// Atomic Operations
INCR(key, delta) -> new_value
CAS(key, expected_value, new_value) -> OK | CONFLICT

Core Entity

Cache Entry
NameDescription
key (PK)Unique identifier for the cached value. Typically a string.
valueThe cached data. Can be serialized objects, strings, or binary.
ttlTime-to-live in seconds. NULL means no expiration.
expiration_timeAbsolute timestamp when the entry expires.
versionMonotonic version for CAS operations and conflict detection.
access_timeLast access timestamp for LRU eviction.
access_countAccess frequency counter for LFU eviction.
created_atEntry creation timestamp for debugging and metrics.

Sign in to continue reading

"Distributed Cache" requires a free account to access.

Sign in to continue