System Design Ch.6: Cache for 1B+ Requests

Outline

Transcript

0:00 Picture this. A massive celebrity changes their profile picture at exactly noon. Oh man, the classic nightmare. Right. But at 12.01, millions of requests are just flooding your system, all asking for that exact same image object. Yep. Now you were a senior backend engineer, you did your homework, you moved that read behind a distributed cache like months ago. So on paper, your database should be perfectly fine. Should be sitting there doing almost nothing. Exactly. But the incident channel still lights up.

0:28 The service starts dropping requests. Why? Well, because that single cache key lives on one single node. And that specific node pegs a CPU core at 100%. And suddenly, your highly available distributed architecture goes completely down. It is, it's really the ultimate architectural wake-up call. You realize in that exact moment that a distributed cache is not just some giant pool of shared memory. Right. It is a complex, heavily contested routing system. I mean, it's an eviction system fighting memory limits, and it's a failure containment system that has to survive topology changes without, you know, crushing your primary data store.

1:05 Which is exactly what we are getting into today. But today, we are bridging the gap. We're moving from just knowing that caches exist to building a really crisp production-grade mental model for how a cache actually behaves as a clustered fleet under extreme load. Which means moving way beyond the beginner's mindset. Right. Because the scale changes everything. It does. At a small scale, a cache is just one instance sitting in front of your database. But at a massive scale, the hard part is no longer deciding whether to cache an item.

1:36 Yeah. The real engineering challenge is knowing where that value lives across, say, a thousand-node cluster. And understanding the actual mathematical mechanics of what breaks when a node dies and the topology shifts. So let's start with how that single node melted down during the celebrity profile update. Because to understand that failure, we first have to look at how a key gets mathematically assigned to a specific server in the first place. Yeah. The routing trap. Right. And the naive approach almost every engineer reaches for first is modular hashing.

2:08 You take the string of the cache key, run it through a hash function to get an integer, and you take modulo n, where n is the total number of cache servers. I mean, the appeal there is totally obvious. It's cheap. It's a decentralized constant-time operation. Right. The client application can compute the placement completely independently without making a network call to some centralized routing proxy. Which sounds great on paper. But the trap springs the moment the cluster scales. So the moment your cluster changes, say you scale from 10 nodes to 11, your n becomes 11.

2:41 And modulo 10 and modulo 11 yield completely different remainders for almost every single integer. Exactly. Suddenly, everyone has a completely different address. Because clients are routing requests to the wrong nodes, every single request registers as a cache miss. Which means they all fall back to the database simultaneously. Right. And that resulting thundering herd completely saturates your database connection pool and just takes the site offline. This is exactly why the industry relies on consistent hashing.

3:09 Okay, let's break that down. Instead of numbering nodes sequentially from 0 to n minus 1, we place both the cache keys and the physical cache nodes onto a conceptual hash ring. A hash ring, right. Yeah. We map the entire output space of our hash function onto a circle. When a request comes in, we hash the key to find its coordinate on that ring. And then we simply walk clockwise until we intersect with the first server node. The mechanics there solve the topology change problem. When a new node joins the cluster, it claims a coordinate on that ring.

3:41 It only assumes responsibility for the keys that fall on the arc between itself and the preceding node. So only that specific slice of keys needs to be repopulated. Yes. The rest of the key space, which could be 90% of your data, stays securely anchored to the original nodes. Consistent hashing buys you that locality of churn. But placing physical nodes on a hash ring introduces a bizarre new problem, doesn't it? Statistical variance. Because if each physical machine is placed onto the ring exactly once, the spacing between them won't be perfectly even.

4:13 Not at all. One node might own a massive 50-degree arc of the ring, absorbing a huge percentage of traffic, while another node lands right next to a neighbor and ends up owning a tiny, like, 2-degree sliver. And that variance is fatal if you have heterogeneous hardware. If you've got older, smaller instances mixed with high-memory instances, you cannot let random hash placement dictate their load. Which brings us to virtual nodes. Yes. Virtual nodes. So virtual nodes are basically like giving each warehouse many small delivery zones scattered randomly all over a city, instead of assigning them one giant continuous district.

4:50 Right. To achieve this mechanically, we don't just hash the node's IP address once. We append an index number to it. We hash node A01, node A02, all the way up to, say, node A256. Placing hundreds of virtual markers for a single physical server smooths out the statistical variance beautifully. It spreads the risk out. And it provides a brilliant mechanism for hardware weighting. Yeah. If you rack a massive, high-capacity server, you simply configure it to generate 500 virtual nodes instead of 200.

5:21 Oh, wow. So it effortlessly absorbs a proportionally larger share of the hash ring? Yep. It's super elegant. Okay. So we've mathematically solved placement. Keys are distributed evenly across our virtual nodes. The thundering herd during a cluster expansion is mitigated. But the next bottleneck we hit in production is physical RAM. Always. You can route perfectly, but every cache in the world is fundamentally smaller than the working set of data it really wishes it could hold. But wait, doesn't TTL handle this?

5:50 The data just expires, right? It frees up memory and everything balances out. That assumption right there is one of the most common causes of degraded cache performance. You must firmly separate the concept of freshness from the concept of capacity. Freshness versus capacity. Yeah. TTL dictates freshness. It answers the question, how long is it safe to trust this cached value before we are legally or logically obligated to query the database again? Right. Because data gets stale. Exactly. But eviction dictates capacity.

6:21 Eviction answers the question, the server's RAM just hit 99%. What specific bytes do we delete right now to make room for an incoming write? I see. So if you rely purely on TTL, your cache just ends up hoarding cold, useless data simply because the clock says it hasn't technically expired yet. Yes. Meanwhile, you might be rejecting or randomly dropping highly valuable, constantly requested data just because the memory is saturated. Exactly. You need a dedicated eviction policy. Yeah. Now the industry standard is LRU, least recently used.

6:53 But I want to pause here because doing strict LRU in a distributed cluster seems mathematically impossible without crippling the system's latency. It is totally impossible. Strict global LRU would require coordinating absolute timestamps across a thousand nodes for every single read operation. That sounds like a nightmare. It is. The network overhead and distributed locking would completely destroy your throughput. So in practice, nodes manage memory locally. Okay. That makes sense. When a specific physical node fills up, it evicts its own least recently used keys independently.

7:26 Furthermore, engines like Redis don't even do strict local LRU. Oh, really? No, they do approximate LRU. They randomly sample a small handful of keys and evict the oldest one from that sample. It saves CPU cycles while achieving like 99% of the practical benefit. Okay. So what about LFU, least frequently used? Because tracking the total historical number of times an item was requested feels like a smarter heuristic than just looking at the last microsecond it was touched. LFU sounds great in a textbook.

7:56 But it completely fails against temporal spikes. Recency is simply a better predictor of immediate future reuse for standard web traffic. The underlying theme here is that average metrics hide catastrophic localized failures. Yeah. We have perfectly routed keys. We're managing local memory beautifully with approximate LRU and admission policies. And yet the system can still crash. Why? Because consistent hashing only solves the average distribution of keys across the cluster. It does absolutely nothing to solve the traffic concentration on one specific key.

8:29 Which brings us right back to the hook. That celebrity profile incident. A hot key. Consistent hashing mathematically guarantees that the celebrity's profile key is routed to one specific node on the ring. Yep. It doesn't matter if you have a thousand massive servers in the fleet. That one key maps to one virtual node which maps to one physical machine. And the CPU on that single machine becomes your system's total ceiling for three -put. So a hot key is basically like hosting an incredibly exclusive celebrity meet and greet, but assigning it to one single entrance door at the stadium.

9:01 Exactly. The stadium itself could seat 100,000 people. It has hundreds of doors. But because everyone in the city wants to see that one celebrity, that one door is absolutely swarmed. The rest of the stadium is empty, but the structural integrity of that one door just fails. Right. So how do we split traffic for an item that logically belongs in only one place on the hash ring? Well, we have to break the rule that it only lives in one place. The standard mitigation for a hot key is replication. You append a randomized suffix to the key name.

9:30 Oh, okay. Instead of writing the profile object to user profile 123, you write it to user profile 123 underscore one underscore two underscore three all the way to say underscore 10. Then, when the client application wants to read the profile, it randomly picks an integer between one and 10, appends it to the key, and makes the request. You have instantly distributed the read volume across 10 different nodes on the hash ring. That's clever. I've also seen teams use a two-level cache for this. They keep a small, in-process local memory cache, like Guava and Java, right inside the application server.

10:04 Yeah. Very common. Sitting directly in front of the distributed network cache. That local cache is incredibly fast. It doesn't even cross the network. And it absorbs the repeated reads from that specific application instance before they ever hit the distributed cluster. Two-level caches are super powerful, but they do heavily inflate your memory overhead across the application layer. True. And more importantly, whether you use key suffixing or local application caches, you've introduced replication.

10:32 And the moment you replicate a single piece of data across multiple isolated memory spaces, you inherit a severe coherence problem. Right. Split brain at the caching layer. Because if that celebrity updates their profile picture again, and you have copies sitting on 10 different cache nodes and 50 different application servers, how do you invalidate them? It gets messy. Do we need strongly consistent cache-to-cache replication? Absolutely not. Attempting strongly consistent cache replication across a distributed fleet is a massive architectural anti-pattern.

11:04 Why is that? Because you end up paying the massive latency costs of distributed consensus protocols, like Paxos or Raft. Just to keep a non-authoritative acceleration layer synchronized. The database is your absolute source of truth. The cache is allowed to be briefly stale. Okay. So what are the actual mechanics of invalidation in that scenario? How do we tear down those 50 copies without grinding the whole system to a halt? You decouple the invalidation from the write path. A common pattern is having a background process like Debezium.

11:36 Read the transaction log directly from the primary database. It publishes an event to a message broker like Kafka. A fleet of consumer workers listens to that topic and aggressively issues delete commands to all known cache nodes containing that key. Oh, I see. It's an eventually consistent purge. Some clients might see the old profile picture for a few hundred milliseconds, but that is the tradeoff you accepted the moment you decided to cache. Let the next organic read miss against the empty cache and lazily repopulate the new value from the database.

12:06 Okay, so replicating hotkeys saves your single node from a CPU meltdown, but it introduces a terrifying new fragility. Because when you deploy a new cluster or execute a regional failover or do a rolling restart, you don't just have one empty node. No, you have a fleet. You have an entire fleet of completely empty replicated nodes. What happens when a huge population of nodes is suddenly empty? How do we prevent that cold start from behaving like a self-inflicted DDoS attack? It's dangerous. If every cache miss falls straight through to the database simultaneously during a failover, your recovery path actually becomes the outage.

12:42 Your connection pools exhaust instantly. The database just locks up. The danger of the miss storm. Exactly. So we need cache warming. And cache warming is essentially like stocking the front display table of a retail store before you unlock the doors for a major holiday sale. Yes. If you wait for the massive rush of customers to flood in and demand products to fill those empty shelves dynamically, the checkout line essentially becomes your inventory routing system. The entire store halts. You have to put the highest velocity items on the table before you flip the sign to open.

13:15 Exactly. You do not let live traffic warm a cold cluster. You implement background asynchronous workers that replay recent read logs. Or you run scripts that specifically query your known hottest entities from the database at a strictly rate-limited pace. You prefetch them. Right. You prefetch and insert those values into the new cache cluster until its hit rate climbs to a safe threshold. Only then do you update the DNS or the load balancer to rule live client traffic to the new fleet. But this brings up a massive architectural dilemma about routing.

13:47 Who is actually computing the key placement? I mean, who holds the hash ring? That is the million-dollar question. If we use client-side routing, meaning the application microservices hash the key themselves and connect directly to the specific cache node, aren't we creating a network nightmare? Oh, it can be. If I have 500 microservices instances and a cache cluster scales up or down, broadcasting that topology state update to 500 clients simultaneously seems incredibly brittle. Not to mention, those 500 clients each maintaining persistent TCP connections to 50 different cache nodes is, what, 25,000 open connections?

14:25 You have correctly identified the constraints of the smart client approach. Client-side routing is insanely fast because it eliminates a network hop. Sure. But as your microservices architecture scales, the connection overhead and the topology synchronization become overwhelming liabilities. If the clients learn about a node failure slowly, they blindly fire requests into a black hole, generating these really painful timeouts. So we add a proxy layer. Yes. You deploy a fleet of routing proxies, like Envoy or TwemProxy, directly in front of the cache cluster.

14:55 The application clients become totally dumb. They just ask for a key. Exactly. They just open a single connection to the proxy and ask for it. The proxy layer centralizes all the topology knowledge, manages the consistent hash ring, and maintains the massive connection pools to the back-end cache nodes. So the trade-off is basically physics. Yep. You have added an extra network hop, which has a fraction of a millisecond of latency, and you've introduced a new potential bottleneck that needs its own high -availability setup.

15:22 This is one of those classic design decisions that looks like simple plumbing on a whiteboard until a massive incident under peak load proves it was the defining architectural constraint. And speaking of incidents, you cannot manage this failure architecture if you are looking at the wrong dashboards. Oh, observability is critical. And engineers get this wrong all the time. A 95% average hit rate looks phenomenal on a dashboard. Engineers will literally high-five over it. I've given those high-fives.

15:49 Right. But that average metric completely masks localized disaster. What if the remaining 5% of misses are the most complex multi-table JOIN queries your database executes? Oh, man. Or what if the hit rate is 95% globally? But because of a hot key, one specific physical node is thrashing at 100% CPU, while the rest of the fleet is completely idle. The system is failing, the users are seeing errors, but the dashboard is green. So you have to measure the pain, not just the success. You need per-node CPU utilization, network saturation, and you really need to monitor the actual database load caused specifically by cache misses.

16:31 Precisely. Notice something structural about this entire architectural breakdown, by the way. What's that? We are at the end of the deep dive, and only now, after designing the consistent hash ring, the eviction strategies, the hot key mitigations, the warming pipelines, and the routing proxy, only now do we talk about the brand names. Technology selection always comes last. It really does. But since we are here, let's briefly contrast the two giants, Memcached and Redis. Sure. Memcached is the simpler philosophy.

16:58 It is a pure, distributed, multi-threaded, in-memory object cache. Wow. It has a very narrow job description. It stores strings and blobs, and it uses a highly efficient slab allocator designed specifically to prevent memory fragmentation under heavy caching workloads. And Redis is fundamentally a different beast. Redis is a massive data structure server. Yes, it acts as an incredible cache, but it operates on a single-threaded event loop for its core command execution. It gives you rich primitives, sorted sets for leaderboards, atomic counters for rate limiting, geospatial indexes, and server-side Lua scripting.

17:34 Both are phenomenal tools. If you want a straightforward, ephemeral, key-value cache, the multi-threaded simplicity of Memcache is an incredible advantage. But if your layer requires richer operations and complex data mutations without pulling the payload back to the application layer, Redis is obviously the pragmatic choice. Absolutely. But the core lesson of today remains, a badly placed cache cluster built on the right technology is still a fragile, poorly designed system. A badly designed cache is just a distributed system you are forced to rescue during an active incident.

18:06 A well-designed cache is a story about bounded, controlled failure. I love that. Bounded, controlled failure. Well, that wraps up Chapter 6 on distributed caching topologies. Next time, in Chapter 7, we are shifting gears entirely. We are moving away from caching static hot objects and diving into live, stateful connections. That's going to be a big one. It is. We will be designing a massive-scale chat system, unpacking the complexities of WebSockets, presence protocols, and the absolute headache that is guaranteed message ordering across a distributed fanout.

18:39 Keep designing, keep building, and we will see you in the next deep dive.