M3 Design Distributed Cache
Loading learning experience...
Lecture transcript
Read the narration for M3: Design Distributed Cache
Design Distributed Cache
Welcome everyone, today we will design a distributed cache, exploring topology choices, consistent hashing, leases, and how to handle eviction and cross region invalidation.
From WhatsApp to caching: why the read path matters
Dr. Wei: Today is about why the read path becomes the real bottleneck in large systems, even when writes look straightforward. When a user opens a chat, the app often needs to read many small pieces of data quickly, not just one big file.
Dr. Wei: A concrete example: media might be stored as an encrypted blob, and the client receives a URL plus a decryption key. But to render the screen, the client also needs profile info, group membership, message metadata, and thumbnails, which can mean dozens of reads in a burst.
Sam: So the core issue is that one screen load can translate into lots of small reads, and at scale those bursts are what stress the system first, even if the big media file is served elsewhere.
Why caching at Meta scale is non-negotiable
Dr. Wei: Before we talk about algorithms or components, let’s ground ourselves in why caching is not optional at Meta scale.
Sam: When you say non negotiable, is the argument basically that each database shard tops out at around ten thousand reads per second, so even a modest spike forces you to rely on cache for the vast majority of reads?
Dr. Wei: So the practical goal is simple: make the cache handle essentially all reads, with a hit rate above ninety nine percent, so the database stays stable, predictable, and frankly boring.
Requirements: turn "fast" into numbers
Dr. Wei: Before we design anything, we need to turn words like fast and scalable into specific, testable requirements.
Sam: What are the top numbers you want me to state in an interview for a cache, like latency targets, value sizes, and what we assume about consistency?
Dr. Wei: Start by defining the traffic: let R day be the total reads per day for the whole system, and let Q be the average reads per second. Then we can estimate Q by dividing R day by eighty six thousand four hundred seconds in a day. For example, if R day is eight hundred sixty four million reads per day, Q is about ten thousand reads per second on average. We also need to be explicit about consistency: whether eventual consistency is acceptable, and how stale data is allowed to get before time to live plus invalidation bring it back in line.
Cache topologies: who owns the miss path?
Dr. Wei: On this slide, we’re comparing cache topologies by a simple question: when data is not in cache, who is responsible for going to the database and getting it?
Sam: Is it fair to summarize it as: look aside puts the miss handling in the application, while look through moves that logic into the cache service, and write behind is a separate decision about where writes land first?
Dr. Wei: If the cache owns the miss path, the cache service itself can load missing data from the database on demand, and related choices like write-behind add another dimension: who is responsible for pushing writes to the database, and when that happens.
Meta-style look-aside: the critical read path
Dr. Wei: On this slide we will focus on the critical read path in a look-aside cache, meaning the cache is checked first but the cache is not the source of truth.
Sam: In practice, does look aside mean every client team owns the miss behavior and retry behavior, so the main risk is they accidentally create a miss storm that turns into a database spike?
Dr. Wei: We also want to be alert to where race conditions and overload can appear, because the miss path is owned by the client and can amplify traffic to the database if many clients miss at the same time.
Partitioning: consistent hashing buys graceful change
Dr. Wei: When you build a distributed cache, you need a way to split the key space across many servers so reads and writes spread out instead of piling onto one machine.
Sam: What I want to sanity check is the failure behavior: if we lose one cache node out of N, consistent hashing means only about one over N of keys remap, instead of almost everything changing like modulo hashing, right?
Dr. Wei: Consistent hashing is popular because it aims for graceful change: only a small fraction of keys should move when membership changes. In expectation, about one over N of keys move, proportional to the failed node’s share, and using virtual nodes reduces variance so the load and the remapped fraction are less skewed; by contrast, naive modulo hashing tends to remap almost all keys when N changes.
Eviction and cache pools: making memory predictable
Dr. Wei: When a cache is small compared to demand, the hard part is not storing data, it is deciding what gets to stay. If that decision is unpredictable, your latency and hit rate will swing wildly as traffic patterns change.
Sam: Is the practical reason for cache pools and slab style memory management to keep one noisy workload from evicting another workload’s hot keys, and to avoid wasting memory on fragmentation?
Dr. Wei: So the goal of this section is to make memory behavior predictable: choose an eviction policy that matches your access pattern, organize memory to avoid waste, and isolate different kinds of traffic so one group cannot erase the benefit for everyone else.
The thundering herd: one hot key takes down the database
Dr. Wei: Let’s understand the failure mode before we talk about fixes. We have a single hot cache key that many users depend on, and it expires around the same moment. When that happens, the cache stops absorbing load and the system behavior can change suddenly and dramatically.
Dr. Wei: Before we reveal what happens next, Sam, make a quick prediction: if about one hundred thousand requests arrive for that same profile right as the key expires, what kind of read pressure does the database see, and what is one cascading mechanism you expect to kick in first, like connection pool exhaustion, queue buildup, or retries?
Sam: So if one hundred thousand requests show up and they all miss, the database effectively sees one hundred thousand reads for the same item, compressed into seconds, which is an instant overload for most shards.
Dr. Wei: Exactly. The outcome is that they all fall through to the database, and it’s not just raw query rate. Connection pools and request queues amplify the failure: threads block waiting for database connections, queues fill up, timeouts trigger retries, and those retries add even more pressure. So a single expiring hot key can turn into a sudden wave of database traffic that overloads the database and then cascades into broader service failures.
Leases: coordinate misses without central locks
Dr. Wei: When multiple clients request the same missing key at the same time, we need a clear rule for who is allowed to fetch from the database and how everyone else should wait and retry.
Sam: I have seen single flight patterns in code. Is a lease basically the cache-native version, where only one client is allowed to refill, and everyone else backs off instead of dogpiling the database?
Dr. Wei: Think of this as a small extension to the cache contract: a get can return hit, or miss with lease granted for a short time window, or lease pending which means wait briefly and retry with backoff. The flow is: the first client to see the miss is granted the lease token and becomes the temporary owner for that key. It reads the database once, writes the value back only if it still has the token, and then releases the lease. Meanwhile, other clients do not hammer the database; they see lease pending, sleep for a short time, retry, and eventually get a hit, or the lease expires and a new client can be granted the lease.
Leases also prevent stale sets (the sneaky correctness bug)
Dr. Wei: So far we have focused on performance, but now we need to talk about a sneaky correctness bug: a cache can accidentally bring an old value back to life after a newer value already exists.
Sam: Let me restate it: a slow reader can fetch an older value, then after someone else updates the database, the slow reader blindly sets the cache and effectively resurrects stale data.
Dr. Wei: This happens in a naive look aside cache because the cache set is blind: it accepts the write without knowing whether the writer is using the most recent version. A lease fixes that by forcing a writer to prove it is still authorized to set the value; if the lease is no longer valid, the cache rejects the stale write.
Invalidation: McSqueal makes staleness short-lived
Dr. Wei: When you use a cache in front of a database, the scary moment is right after a write: the database has the new value, but the cache might still have the old one. This slide is about shrinking that staleness window so it is short-lived and predictable.
Sam: So instead of waiting for time to live to eventually expire, we try to actively delete the affected keys right after the database commit, so the staleness window is more like the invalidation propagation time.
Dr. Wei: Two practical details matter. First, mapping a write to the exact cache keys is not automatic in real systems with denormalization and fanout, like a user update affecting multiple views. Second, the event stream is typically at least once, so deletes must be idempotent and safe under duplicates or occasional misses, and you often rely on versioning or careful key design to make this robust.
Failure modes and mitigations (what $IC6$ sounds like)
Dr. Wei: Let’s recap two failure modes that show up in real cache incidents, and pair each one with the simplest guardrail that keeps the rest of the system safe.
Sam: For cache host loss, the part I want to say cleanly is: we will remap a fraction of keys and traffic, so we need spare capacity and throttling to avoid the remaining nodes and the databases seeing a sudden surge.
Dr. Wei: Mitigation: plan headroom and add rate limits so the remapped load is absorbed smoothly. In other words, expect a fraction of keys to move when a node disappears, and keep enough spare capacity that the cluster stays stable.
Dr. Wei: Second: hot key updates. If many clients update or invalidate the same popular key at once, you can create write amplification and constant churn, even when the underlying value barely changes.
Sam: And the mitigation story there is to coalesce or batch the invalidations and updates, so the system does one effective change instead of turning it into thousands of writes per second on the same key.
Exit ticket: size it, then stop the herd
Dr. Wei: For this exit ticket, you will do one quick sizing calculation and one short design reflection about preventing a cache stampede.
Sam: Okay, so with ninety nine point five percent hit rate, the miss rate is zero point five percent. On two million gets per second, that is ten thousand misses per second, which means about ten thousand reads per second to the database just from misses.
Dr. Wei: Second, choose one approach to reduce the herd effect when a key expires: leases, or TTL-only. Then defend your choice with a specific failure story, like what happens during a hot key spike, a deploy, or a brief database slowdown.
Thank you for watching!
Thanks for watching. Subscribe and share if you found this useful—see you next time!