Distributed Caching: Building Fast and Resilient Cloud Applications

Modern cloud applications frequently read the same data many times. Repeatedly fetching that data from a database can increase latency, consume database connections, and create unnecessary load on the primary data store.

Distributed caching places frequently accessed data in a low-latency memory-based system that can be shared across multiple application instances. Instead of every request reaching the database, applications can retrieve commonly used data directly from the cache.

Caching is useful for user profiles, product catalogs, configuration, session data, authorization information, API responses, computed results, and other data that can tolerate controlled freshness.

The central engineering challenge is deciding what should be cached, how long it should remain valid, how stale data should be handled, and how the system should behave when the cache becomes unavailable.

Caching Models: Cache-Aside, Read-Through, and Write-Through

Different caching strategies determine how applications interact with the cache and the underlying database.

1. Cache-Aside

  • Mechanism: The application first checks the cache. If the item is missing, it reads from the database and then stores the result in the cache.
  • Strengths: Simple and flexible because the application controls which data is cached.
  • Trade-off: Cache-miss logic must be implemented correctly by application developers.

2. Read-Through Cache

  • Mechanism: The application requests data from the cache while the caching layer automatically loads missing data from the backing store.
  • Strengths: Simplifies application code.
  • Trade-off: Requires caching infrastructure capable of integrating with the underlying data source.

3. Write-Through Cache

  • Mechanism: Writes are sent to the cache and synchronously propagated to the backing data store.
  • Strengths: Keeps cached and persistent data closely synchronized.
  • Trade-off: Write latency can increase because both cache and persistent storage participate in the operation.

4. Write-Behind Cache

  • Mechanism: The application writes to the cache first and persistence happens asynchronously.
  • Strengths: Can provide very fast write responses and absorb write bursts.
  • Trade-off: Data can be lost if cached updates are not persisted before a failure.

Cache Lifecycle: TTL, Expiration, and Invalidation

Cached data must eventually become invalid or be replaced. Otherwise, the cache can return information that no longer reflects the source of truth.

1. Time-To-Live

A TTL specifies how long an item should remain available in the cache before it expires.

2. Explicit Invalidation

Applications can remove or update cached values immediately after the underlying data changes.

3. Versioned Keys

Applications can encode a version into cache keys so new versions naturally bypass older cached values.

4. Stale Data

  • Short TTLs reduce staleness but increase database traffic.
  • Long TTLs improve cache hit rates but increase freshness risk.
  • Critical data may require explicit invalidation.
  • Non-critical data can often tolerate eventual consistency.

Cache Eviction Policies

Memory is limited, so a distributed cache must decide which entries to remove when available capacity is exhausted.

1. Least Recently Used (LRU)

Removes entries that have not been accessed recently.

2. Least Frequently Used (LFU)

Removes entries with the lowest access frequency.

3. First-In, First-Out (FIFO)

Removes entries based on insertion order.

4. TTL-Based Eviction

Entries are removed after their configured expiration time.

  • LRU works well for workloads with strong temporal locality.
  • LFU can protect frequently accessed hot data.
  • TTL should reflect how long data can safely remain stale.
  • Eviction policies should be selected based on actual workload behavior.

Cache Consistency and Stale Data

Caching introduces a second copy of application data, creating consistency challenges between the cache and the source of truth.

1. Stronger Consistency

Applications explicitly invalidate or update cached data whenever the source data changes.

2. Eventual Consistency

Cached data may temporarily differ from the database but converges after expiration, invalidation, or asynchronous updates.

3. Stale-While-Revalidate

The cache can temporarily return slightly stale data while asynchronously refreshing the value in the background.

  • Do not cache data without understanding its freshness requirements.
  • Separate strongly consistent workflows from eventually consistent reads.
  • Use explicit invalidation for data where stale values can cause serious business errors.
  • Document acceptable staleness for each cached dataset.

Cache Stampede, Thundering Herd, and Hot Keys

A cache can reduce backend load dramatically, but poorly designed expiration behavior can cause sudden traffic spikes against the database.

1. Cache Stampede

A popular cache entry expires and thousands of application requests simultaneously attempt to rebuild the same value.

2. Request Coalescing

Only one request rebuilds a missing value while other requests wait for the shared result.

3. Probabilistic Early Refresh

Applications can refresh popular entries slightly before expiration to reduce synchronized cache misses.

4. Hot Keys

A small number of extremely popular keys can overload a single cache node or create concentrated backend traffic.

  • Use jittered expiration times.
  • Use request coalescing for expensive cache misses.
  • Replicate extremely hot values when necessary.
  • Monitor cache-miss spikes and key-level access patterns.

Distributed Cache Sharding and Partitioning

A single cache node cannot provide unlimited memory or throughput. Large systems distribute cache keys across multiple nodes.

1. Hash-Based Sharding

A hash function maps each cache key to a specific cache node.

2. Consistent Hashing

Consistent hashing reduces the number of keys that need to move when cache nodes are added or removed.

3. Virtual Nodes

Multiple logical positions can represent each physical cache node, improving distribution across the hash ring.

4. Rebalancing

  • Adding nodes should distribute workload without causing massive cache invalidation.
  • Removing nodes should gracefully redistribute keys.
  • Avoid concentrating popular keys on a single node.
  • Monitor memory utilization and request distribution per node.

Cache Replication and High Availability

Distributed caches can replicate data across multiple nodes to improve availability and tolerate individual cache failures.

1. Primary-Replica Architecture

A primary node handles writes while replicas maintain copies of cached data.

2. Automatic Failover

If a primary cache node fails, another node can take over responsibility.

3. Replication Lag

Asynchronous replication can temporarily leave replicas behind the latest primary state.

  • Cache replicas improve availability but do not automatically guarantee strong consistency.
  • Applications should tolerate cache loss because the persistent database remains the source of truth in many architectures.
  • Failover behavior should be tested before production incidents occur.

Distributed Locks and Cache Coordination

Some caching workloads require coordination so multiple application instances do not perform the same expensive operation simultaneously.

1. Cache Rebuild Lock

A temporary lock can ensure that only one worker rebuilds an expensive missing cache entry.

2. Lock Expiration

Locks should have bounded expiration times so a crashed worker does not permanently prevent other workers from making progress.

3. Lock Contention

Highly popular keys can create excessive contention if every request waits on the same lock.

  • Keep lock duration short.
  • Always use an expiration mechanism.
  • Avoid using distributed locks when simpler request coalescing is sufficient.
  • Do not assume a cache lock automatically provides transactional guarantees for database operations.

Cache Failure and Graceful Degradation

A cache should generally improve application performance rather than become a mandatory single point of failure.

1. Cache Unavailable

Applications can fall back to the database when the cache cannot be reached, provided the database has sufficient capacity.

2. Partial Cache Failure

If one cache node fails, sharded systems should continue serving unaffected keys while rebuilding or redistributing the failed partition.

3. Database Protection

A complete cache outage can suddenly multiply database traffic. Systems should use connection limits, request throttling, and controlled fallback behavior to prevent a database overload.

  • Set cache operation timeouts.
  • Avoid waiting indefinitely for cache responses.
  • Protect the database from cache-miss storms.
  • Use fallback data where appropriate.
  • Monitor cache health independently from application health.

C++ Conceptual Simulation Blueprint (Cache-Aside Pattern)

C++
Example conceptual distributed cache-aside implementation
#include <iostream>
#include <string>
#include <unordered_map>

class Cache {
private:
    std::unordered_map<std::string, std::string> values;

public:
    bool get(const std::string& key, std::string& value) {
        auto it = values.find(key);

        if (it == values.end()) {
            return false;
        }

        value = it->second;
        return true;
    }

    void set(const std::string& key,
             const std::string& value) {
        values[key] = value;
    }
};

class Database {
public:
    std::string read(const std::string& key) {
        std::cout << "Database read: "
                  << key << std::endl;

        return "value-for-" + key;
    }
};

class Application {
private:
    Cache& cache;
    Database& database;

public:
    Application(Cache& c, Database& d)
        : cache(c), database(d) {}

    std::string get(const std::string& key) {
        std::string value;

        if (cache.get(key, value)) {
            std::cout << "Cache hit: "
                      << key << std::endl;
            return value;
        }

        std::cout << "Cache miss: "
                  << key << std::endl;

        value = database.read(key);
        cache.set(key, value);

        return value;
    }
};

Distributed Cache Performance and Observability

Cache performance should be measured using both cache-level metrics and the effect caching has on downstream systems.

  • Cache Hit Rate: Percentage of requests successfully served from the cache.
  • Cache Miss Rate: Percentage of requests requiring a backend lookup.
  • Cache Latency: Time required to retrieve or store cached values.
  • Eviction Rate: Number of entries removed because of memory pressure or eviction policy.
  • Memory Utilization: Percentage of cache memory currently occupied.
  • Hot Key Frequency: Number of requests concentrated on the most frequently accessed keys.
  • Backend Load: Database requests generated by cache misses.
  • Stampede Events: Number of coordinated or simultaneous cache rebuilds.
  • Replication Lag: Difference between primary and replica cache state where applicable.
  • Cache Availability: Percentage of time the cache successfully handles requests.

Real-World Cloud & Distributed Caching Implementations

  1. Redis: In-memory data platform commonly used for distributed caching, counters, sessions, queues, locks, and fast application state.
  2. Memcached: Lightweight distributed memory caching system commonly used to cache frequently accessed application data.
  3. Amazon ElastiCache: Managed AWS caching service supporting Redis-compatible and Memcached-based workloads.
  4. Amazon DynamoDB Accelerator (DAX): Managed in-memory caching layer designed to accelerate DynamoDB workloads.
  5. Google Cloud Memorystore: Managed in-memory caching service supporting Redis and Memcached workloads.
  6. Azure Managed Redis: Managed Redis-based caching infrastructure for applications running on Azure.
  7. CDN Edge Caching: Static and cacheable HTTP content can be stored near users to reduce latency and origin-server traffic.
  8. Microservice Caching: Individual services can cache frequently accessed database records or expensive computed results to reduce dependency load.