Skip to content

The finished design, decision by decision

How to design a Distributed Cache (Memcache)

Not the one correct diagram, but a design you can defend under these constraints: the finished architecture, then every stage's question with the reasoning that answers it, the tradeoffs it accepts, and where another engineer could land differently.

A look-aside cache at Facebook's scale

The short answer

7 parts, each with one job. The map below shows how requests and data move between them; the stages after it explain why each part is there.

Users
Load pages.
Web servers
Render pages; read through the cache, query MySQL on a miss, delete keys after writes.
mcrouter
Routes each key to a memcached server by consistent hashing; fans out deletes.
memcached pool
Demand-filled key-value cache; issues leases on misses.
Gutter pool
About 1% of servers, idle until a memcached server fails; short-lived entries.
MySQL
The source of truth. Writes embed the cache keys to invalidate.
Invalidation daemon
Tails the commit log, extracts deletes and sends them in batches to every cluster.
1234567CLIENTUsersSERVICEWeb serversSERVICEmcrouterCACHEmemcached poolCACHEGutter poolDATABASEMySQLWORKERInvalidationdaemon

Select a component to see what it is responsible for and which state it owns.

  1. 1Users → Web servers: Page request
  2. 2Web servers → mcrouter: get / multiget, delete
  3. 3mcrouter → memcached pool: Keys by consistent hash
  4. 4Web servers → MySQL: Query on miss; writes
  5. 5MySQL → Invalidation daemon: Committed deletes
  6. 6Invalidation daemon → mcrouter: Batched deletes
  7. 7mcrouter → Gutter pool: On server failure
  • Request / response
  • Bulk data
  • Asynchronous

Why does this design work?

The cache is a disposable, demand-filled copy: writes go to MySQL and delete cached keys, reads refill on a miss. Because database load equals the miss rate, every mechanism protects the miss path.

Leases reject refills that started before a newer write (no stale sets) and, rate-limited, let one caller refill a hot key while others wait (no thundering herds). A gutter pool absorbs a dead server's load; cold-cluster warmup borrows a neighbour's memory. Invalidations are recorded with each commit and delivered from the log in batches to every cluster, after the data has replicated, so they are durable, replayable and correctly ordered. Writers see their own changes through local deletes and remote markers; everyone else gets bounded staleness.

Invariants, and where they are enforced

  • A value read before a write cannot be stored in the cache after that write's delete.

    A miss returns a lease token; a delete invalidates outstanding tokens; a set with an invalid token is rejected. Enforced by memcached pool.

  • Every committed write's invalidations reach every cluster, and can be replayed if lost.

    Keys to delete are recorded with the committed SQL; daemons read the commit log and batch deletes to each cluster. Enforced by MySQL, Invalidation daemon.

  • A failed cache server does not send its full load to the database.

    Requests to an unresponsive server retry against a small gutter pool that fills on demand and expires quickly. Enforced by mcrouter, Gutter pool.

What does it rely on?

  • Reads vastly outnumber writes, and most data tolerates brief staleness.
  • The cache is never the only copy of anything (except deliberately, like remote markers).
  • Cache keys to invalidate can be known at write time.
  • Spare capacity (gutter) is kept idle for failures.

What tradeoffs does it make?

ChoiceGainsCosts
Delete on writeIdempotent, order-insensitive invalidation.A miss after every write.
LeasesNo stale sets; one refill per hot key.Brief waits for other readers; a protocol change.
Gutter poolServer failures do not reach the database.~1% idle capacity; slightly stale entries.
Invalidation from the commit logDurable, replayable, batched, correctly ordered.Daemons on every database; invalidation delay.

What are the reasonable alternatives?

Write-through graph cache (like Facebook's later TAO)
Better when the data model is uniform enough for the cache to understand and update it.
Short TTLs with no invalidation
Better when staleness of seconds is fine and write rates are high.
Database read replicas instead of a cache
Better when queries are too varied to cache by key.

When does it stop working?

  • Data requires strong consistency for every reader.
  • Write rates approach read rates, so deletes keep the hit rate low.
  • Keys to invalidate cannot be determined from the write.

Every stage, decided and explained

Spoilers, for the whole investigation: each stage's question and its answer, the reasoning behind it, and the tradeoffs it accepts. If you have not worked through the stages yet, you may want to do that first.

Work through the stages

Stage 1 of 9 · Model

How a look-aside cache behaves

Reads: get from the cache; on a miss, query MySQL and set the result. Writes: update MySQL, then do something about the cached copy.

What you need to know first

A look-aside (or cache-aside) cache sits next to the database, not in front of it. The application does the work:

  • Read: ask the cache. On a hit, done. On a miss, query the database and put the result in the cache.
  • Write: update the database, then deal with the cached copy.

The cache never talks to the database itself. It's a disposable copy the application manages. See Caching.

The cache serves 1,000,000 reads a second with a 99% hit rate. How many reads a second reach the database?

About 10,000 per second.

1% of 1,000,000 = 10,000 a second. The database is provisioned for this miss load, not the total.

The hit rate drops from 99% to 98%. What happens to database load?

It roughly doubles: misses go from 1% to 2%.

Database load follows the miss rate, not the hit rate. A one-point drop in hits is a 100% increase in misses.

On write, there are two options for the cached copy: set the new value, or delete the key and let the next read refill it.

Deletes are idempotent (doing it twice is the same as once) and order-insensitive (two deletes in either order leave the same result). Two sets racing can arrive in the wrong order and leave the older value cached.

What the stage asks

Which statements hold?

  1. Fails

    After a write, the web server should set the new value in the cache so the next read hits.

    Facebook deletes instead. A delete is idempotent and order-insensitive: two deletes in any order leave the same result. Two concurrent sets can arrive in the wrong order and leave the older value cached. The next read refills from the database.

  2. Holds

    Because the cache is not the source of truth, losing any cached key must always be safe.

    It may cost a database query, but never correctness. Hold on to this: stage 7 introduces a key for which it is no longer true.

  3. Fails

    To keep latency low, a page that needs 500 keys should fetch them one by one.

    500 sequential round trips is far too slow. Clients batch independent keys into parallel multigets, ordering them by data dependencies. That creates its own problem: hundreds of responses arriving at once (incast), which clients limit with a sliding window of outstanding requests.

  4. Holds

    If the hit rate drops from 99% to 98%, the database receives about twice as many reads.

    Misses go from 1% to 2% of reads. Small changes in hit rate are large changes in database load, which is why everything later in this investigation is about protecting the miss path.

The reasoning

  1. In a look-aside cache, the database is the truth, the cache is a disposable copy, and writes delete.
  2. Database load equals the miss rate, so small hit-rate drops are large database load increases.
  3. Deletes are idempotent and order-insensitive; racing sets can leave old values cached.

The look-aside contract is simple: the database is the truth, the cache is a disposable copy, writes delete. The arithmetic is the part people miss: at high hit rates, the database's load is the miss rate, so anything that causes a burst of misses (a hot key, a dead server, an empty cluster) is a database incident.

Stage 2 of 9 · Break it

A stale value that never leaves

Two web servers, A and B, touched the same key around the same time. Find the lines that explain why the cache holds the old value indefinitely.

What you need to know first

A reader that misses does two things at two different times: it reads the database, then later sets the cache. Anything can happen in between, including a write and its delete.

If the write's delete arrives before the reader's set, the delete has nothing to remove, and the set then stores a value read before the write.

The stale value 'Ann' is now in the cache, with no TTL. How long does it stay?

Until the next write to that key or until it's evicted, which could be days

Nothing knows it's stale. The only thing that removes it is another delete or memory pressure.

The fix is a token: on a miss, the cache hands the reader a lease token for that key. A delete for the key invalidates outstanding tokens. The reader's set must carry its token, and the cache rejects sets whose token was invalidated.

The idea is the same as a fencing token: a write is accepted only if nothing newer has happened since it was authorised.

Would a 1-hour TTL on every key fix the stale set?

It bounds the damage to at most an hour of a wrong name, but doesn't prevent the race. Shorter TTLs bound it more tightly, at the cost of more misses for everyone. The lease prevents it outright.

What the stage asks

Select the lines that are part of the problem.

TimelineKey user:42
  1. 1t1 A: get user:42 → miss
  2. 2t2 A: SELECT name FROM users WHERE id = 42 → 'Ann'
  3. 3t3 B: UPDATE users SET name = 'Annie' WHERE id = 42; COMMIT
  4. 4t4 B: delete user:42 (nothing cached yet, so nothing removed)
  5. 5t5 A: set user:42 = 'Ann'

    A's set carries a value read before B's write, and arrives after B's delete. Nothing tells the cache that this value is older than the write it already saw invalidated.

  6. 6t6 memcached: stores 'Ann' (sets are unconditional)

    An unconditional set lets any late writer win. The cache needs a way to reject sets that started before a delete.

  7. 7t7+ everyone: get user:42 → 'Ann' (until the next write to user 42, or eviction)

What the fix has to do

  • The refill read happened before the write but the set arrived after the delete.
  • Nothing removes the stale value until another write or an eviction.
  • Leases: a miss returns a token tied to the key; a delete invalidates it; a set with an invalidated token is rejected.
  • A TTL bounds the damage but does not prevent it.Supporting

The reasoning

  1. A refill read before a write but set after its delete caches the old value with nothing to remove it.
  2. Leases tie a refill to its miss; a delete invalidates the token and the late set is rejected.
  3. TTLs bound staleness but don't prevent the race.

This is a stale set: a refill computed from old data lands after the invalidation meant to remove it. Deleting instead of setting on write does not fix it, because the problem is the reader's set.

Facebook's fix is a lease: a miss hands the client a token for that key, and the refill is accepted only if it carries the token. A delete for that key invalidates outstanding tokens, so A's late set is rejected. It is the same idea as a fencing token in Leases and fencing tokens: a write is only accepted if nothing newer has happened since it was authorised.

Stage 3 of 9 · Decide

The key everyone wants

The database is provisioned for the normal miss rate. This one key is producing a large share of all misses.

What you need to know first

A thundering herd (or stampede) happens when many readers miss the same key at the same moment, and all of them go to the database to rebuild it. With a hot key, that moment is every time the key is deleted or expires.

Raise the request rate and the rebuild time and watch how many queries reach the database with each strategy.

A hot key expires. Rebuilding it takes one expensive query. Until that query finishes, every request for the key misses. The database can run about 200 of these queries at once before it slows down.
5,000/s400 ms

Queries reaching the database while the key is rebuilt (log scale). The dashed line is what it can run at once.

  • No protection2,000
  • Coalesce per server20
  • One rebuild, fleet-wide1
  • Serve stale, refresh behind1
No protection:
Every request that misses queries the database. 2,000 requests wait about 400 ms, longer, because the database is past capacity and every query slows down.
Coalesce per server:
Each of 20 servers lets one request rebuild; the rest on that server wait for it. 2,000 requests wait about 400 ms.
One rebuild, fleet-wide:
A short lock in the cache lets one request in the whole fleet rebuild. 2,000 requests wait about 400 ms.
Serve stale, refresh behind:
Keep serving the old value while one background request rebuilds it. Nobody waits; a few requests see a value that is a moment old.

Without protection, what decides how many queries a hot key's miss sends to the database?

Requests per second for the key × how long the rebuild takes

Every request that arrives during the rebuild misses. At 5,000 a second and 400 ms, that's 2,000 queries for one key.

Request coalescing lets one caller do the rebuild while the others wait for its result. Done inside the cache, it's a lease that is only granted once per interval per key: one reader gets the token, the rest are told to wait a few milliseconds and retry. See Request coalescing.

For data that tolerates being slightly old, waiting isn't even necessary: serve the previous value while one caller refreshes it.

What the stage asks

What do you change?

  1. Flawed

    Give the key a much longer TTL

    The key is not expiring; it is being deleted by writes. A TTL does nothing about misses caused by invalidation.

  2. Sound

    Rate-limit leases: hand out one refill token per key every few seconds; other missing clients are told to wait briefly and retry

    One client refills while the rest wait a few milliseconds and then find the value in the cache. Facebook issued at most one token per key every 10 seconds; on keys prone to herds, peak database queries fell from 17,000 a second to 1,300.

  3. Flawed

    Set the new value on write instead of deleting, so there is never a miss

    It reintroduces racing sets: concurrent updates can land in the wrong order and leave an older count cached. It trades a load problem for a correctness problem.

  4. Defensible

    Keep recently deleted values briefly, and serve them (marked stale) to callers that can tolerate it while one refills

    Facebook did this too, for data where a slightly old value is fine (a like count usually is). It removes the wait entirely, but each caller must opt in, because some data must never be served stale.

What a strong answer covers

  • Every write deletes the key, and many concurrent readers miss in the gap before refill.
  • Only one caller should refill; the others wait or use a stale value.
  • The database sees at most one refill per key per interval, however many readers there are.

The reasoning

  1. A hot key's miss sends (request rate × rebuild time) queries to the database without protection.
  2. Rate-limited leases let one caller refill while others wait briefly: request coalescing in the cache.
  3. Serving a slightly stale value removes the wait for data that tolerates it.

The same lease that prevents stale sets also prevents herds once you rate-limit how often it is granted. That is Request coalescing implemented inside the cache: one caller does the work, everyone else waits on its result.

Stale values are the complementary tool: when a caller can live with a slightly old answer, it does not need to wait at all.

Stage 4 of 9 · Break it

Write the lease-aware read

The cache's get now returns one of three things: a value, a lease token (you should refill), or "wait" (someone else is refilling). setWithLease returns false if the token was invalidated. Write the web server's read helper.

What you need to know first

A lease-aware get returns one of three results, and the client has a branch for each:

ResultMeaningClient does
hitthe valuereturn it
leaseyou should refillload from the database, set with the token
waitsomeone else is refillingsleep briefly, ask again

The client loads the value and its setWithLease returns false (a delete invalidated the token). What should it return to its caller?

The value it loaded: it was correct when read, it just mustn't be cached.

A rejected set only means the cache shouldn't keep this value. Returning it to the current request is fine.

The web server holding the lease crashes before refilling. What happens to the clients told to 'wait'?

Their retries keep getting "wait" until the lease expires on the server. So the client needs a bound: after a few short, growing, jittered waits, it reads the database directly. The page costs a few milliseconds more instead of failing.

What the stage asks

Implement cachedRead.

Reference implementation

const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms));

export async function cachedRead<T>(key: string, load: () => Promise<T>): Promise<T> {
  for (let attempt = 0; attempt < 5; attempt++) {
    const result = await cache.get<T>(key);
    if (result.kind === "hit") return result.value;
    if (result.kind === "lease") {
      const value = await load();
      // Rejected if a delete arrived meanwhile: our value may be stale, so it is
      // not cached, but it is still a valid answer for this request.
      await cache.setWithLease(key, value, result.token);
      return value;
    }
    // Someone else holds the lease; it is usually filled within milliseconds.
    await sleep(5 * 2 ** attempt + Math.random() * 5);
  }
  return load();
}
  • A rejected set is not an error. It only means this value must not be cached; returning it to the caller who asked is fine, because it was correct when read.
  • Waits are short and grow, with jitter; most waiters find the value on their first retry.
  • The final fallback goes to the database, so a lost lease holder (a crashed web server) costs a few milliseconds, not a broken page. Leases also expire on the server for the same reason.

What a strong answer covers

  • Returns the value on a hit.
  • On a lease, loads from the database and sets with the token, returning the loaded value whether or not the set is accepted.
  • On 'wait', sleeps briefly and retries, with a bounded number of attempts.
  • After the retries run out, reads from the database directly rather than failing the page.
  • Adds jitter to the wait so waiting clients do not retry in lockstep.Supporting

The reasoning

  1. Handle all three outcomes of a lease-aware get: hit, refill with the token, or wait.
  2. A rejected set isn't an error; return the loaded value without caching it.
  3. Bound the waiting with jittered retries, then fall back to the database.

The client is where the protocol becomes behaviour: what to do with a token, what to do while waiting, and what to do when waiting goes on too long. Each branch has a reason, and each has a bound.

Stage 5 of 9 · Break it

A cache server dies

The database is provisioned for the normal miss rate. Decide what clients do in the minutes before the replacement arrives.

What you need to know first

Clients choose a cache server for each key by hashing the key. Two common schemes:

  • hash mod N: server = hash(key) % N. Simple, but changing N moves almost every key.
  • Consistent hashing: servers and keys are placed on a ring; each key belongs to the next server clockwise. Adding or removing a server moves only the keys next to it. See Consistent hashing.

Add a node under each placement and compare how many keys move. Then try the ring with one point per node versus a hundred.

2,000 keys spread over cache nodes. Add a node and count the keys that now belong somewhere else: each one is a cache miss, or data to copy.
Placement
4 nodes
4 nodeskeys by hash mod N

Keys per node. The line is a perfectly even share.

  • Node 1500
  • Node 2508
  • Node 3512
  • Node 4480

The busiest node holds 1.0× an even share.

Moving few keys is good for planned changes. A failure raises a different question: the dead server's keys all miss, and their load must go somewhere. Rehashing sends them to the neighbouring servers, which are already busy. And keys aren't equal: one hot key can be a fifth of a server's traffic.

A dead server's hottest key (20% of its traffic) is rehashed onto a healthy server that is already at 85% capacity. What can happen?

The healthy server overloads and fails too, and its keys move on again: a cascade.

Rehashing assumes the load spreads evenly. Hot keys don't spread; they land somewhere whole.

What the stage asks

What should clients do with requests for the failed server's keys?

  1. Flawed

    Go straight to the database for those keys until the server is replaced

    Every request that server used to absorb becomes a database query, minutes of a miss rate the database was never provisioned for. That is how one cache failure becomes a site outage.

  2. Flawed

    Rehash its keys onto the remaining memcached servers

    The failed server's hot keys land on healthy servers that are already busy; a key that is 20% of a server's traffic can overload its new home, which fails in turn. Facebook rejected this specifically because of cascading failures.

  3. Sound

    Retry failed gets against a small, normally idle 'gutter' pool that fills on demand, with short expiry and no invalidations

    Roughly 1% of servers sit idle until needed. A failed get retries against gutter; a gutter miss queries the database once and fills gutter, so the next request hits. Facebook reported gutter hit rates above 35% within four minutes, and a 99% drop in client-visible failures. Entries expire quickly because gutter does not receive invalidations.

  4. Defensible

    Store every key on two memcached servers

    Doubling memory for every key costs a lot to cover rare failures. Facebook did replicate some key families inside pools, for read throughput, not as the general answer to failure.

What a strong answer covers

  • A failed server's keys all miss at once; the database cannot absorb that load.
  • Rehashing onto busy servers risks cascading failure because of hot keys.
  • Spare capacity that is idle until a failure absorbs the load without disturbing healthy servers.
  • Gutter entries may be slightly stale, so they expire quickly.Supporting

The reasoning

  1. A failed cache server's load has to go somewhere; the database is the worst place.
  2. Rehashing onto busy servers can cascade, because hot keys move whole.
  3. An idle 'gutter' pool absorbs the misses with short-lived entries until the server is replaced.

The general lesson: when part of a cache fails, the load it was absorbing has to go somewhere, and the database is the worst place. Gutter gives that load a home that is empty, cheap and temporary.

Notice why gutter can skip invalidations: its entries live for seconds. Short-lived copies need far less consistency machinery than long-lived ones.

Stage 6 of 9 · Change it

Deletes for every cluster

Web servers currently delete keys in their own cluster after writing. Deletes must now reach every cluster, reliably, at a very high rate.

What you need to know first

With several clusters, a popular key may be cached in each of them, and a write must reach every copy. Two properties matter for that delivery:

  • Durable: a delete must not be lost because the web server that wrote crashed.
  • Efficient: deletes are frequent, so they need batching across cluster boundaries.

A web server commits a write, then crashes before sending deletes to the other clusters. What do those clusters keep serving?

The old value, indefinitely: nothing else knows the key changed. The write is durable in the database, but the obligation to invalidate lived only in the crashed process's memory.

The fix records the invalidation with the write: the keys to delete are part of the committed transaction, and daemons that tail the database's commit log send them out. The log is durable and replayable, so a crash or a delivery bug can be recovered from. It's the Transactional outbox idea applied to caches.

Deletes now flow from the commit log. Why does the writing web server still delete the key in its own cluster right away?

So the user who wrote sees their change on their next request, without waiting for the log pipeline

The log path has some delay. The local delete gives the writer read-your-writes immediately.

What the stage asks

How should invalidations reach every cluster?

  1. Defensible

    The web server that wrote sends deletes to every cluster

    Simple, but each web server batches poorly, so packet rates across cluster boundaries are high. And if deletes are lost or misrouted (a configuration bug), there is no record to replay; Facebook's recourse used to be a slow rolling restart of the whole cache.

  2. Sound

    Record the keys to invalidate with the committed SQL; daemons on each database tail the commit log and send batched deletes to routers in every cluster

    Invalidations become part of the durable commit, so they cannot be lost when a web server crashes, and they can be replayed. Daemons batch them (an 18× improvement in deletes per packet), and routers in each cluster fan them out. The writing web server still deletes in its own cluster for read-your-writes.

  3. Defensible

    Rely on short TTLs everywhere instead of invalidating

    Staleness is bounded with no machinery, but short TTLs mean many more misses across the whole fleet, which is exactly the load the cache exists to remove.

  4. Flawed

    Write the new value through to every cluster's cache

    Sets racing across clusters can arrive out of order and leave old values cached, and most of those clusters may never read the key, so you fill caches with data nobody asked for.

What a strong answer covers

  • Invalidations committed with the write cannot be lost when the writer crashes.
  • A log can be replayed after a delivery failure or misrouting.
  • Daemons batch deletes efficiently across cluster boundaries.
  • The writer still deletes locally for read-your-writes.Supporting

The reasoning

  1. Record invalidations with the commit and deliver them from the log: durable, replayable, batched.
  2. The writer still deletes locally for read-your-writes.
  3. Most deletes hit nothing, so cheap batched delivery matters more than precise targeting.

This is the Transactional outbox idea applied to caches: the side effect (invalidate these keys) is recorded in the same commit as the data change, and a separate process reliably delivers it from the log. See also Append-only logs.

Facebook noted that only about 4% of deletes actually remove a cached value. Most keys are not cached in a given cluster when they change, which is why cheap, batched delivery matters more than precise targeting.

Stage 7 of 9 · Change it

A second region

Replicas can lag behind the master. Caches in the replica region are filled from the local replica.

What you need to know first

With replication, writes go to a primary and are copied to replicas some time later: usually under a second, sometimes much longer. A replica that hasn't received a write yet serves the old data. See Replication.

A cache filled from a lagging replica stores that old data.

The master region writes, then immediately sends a delete to the replica region's cache. The data arrives at the replica 2 seconds later. A read in the replica region happens in between. What gets cached?

The read misses (the key was deleted), refills from the replica, which still has the old value, and caches it. The delete has already happened, so nothing removes the old value.

Deletes need to arrive after the data. Tailing each region's own replica log achieves that: the invalidation is applied when the data change is.

For the user who wrote, a remote marker helps: before writing, set a marker key in the local cache saying "recently written". A miss that finds the marker reads from the master region (slower, but current) instead of the local replica.

Why is evicting a remote marker different from evicting a normal cached value?

The marker is information, not a copy: losing it means a read may go to the stale replica.

A normal value can always be rebuilt from the database. A marker's presence is the only record that the replica may be behind.

What the stage asks

Which statements hold?

  1. Holds

    A user in the replica region who changes their profile may see the old version on their next request.

    Their write went to the master. If their next request misses the cache, it refills from a local replica that may not have the change yet, and caches the old value.

  2. Fails

    The master region's web server should send invalidations straight to the replica region right after its write.

    The delete can arrive before replication does. A read in the replica region then misses, refills from the lagging replica, and caches the old value with nothing left to remove it. Invalidations in each region should come from that region's own database log, after the data has arrived.

  3. Holds

    Setting a 'recently written' marker for a key, and sending misses for marked keys to the master region, trades latency for freshness.

    Facebook called these remote markers: set the marker, write to the master, delete the local key. A miss that finds a marker reads from the master region (slower) instead of the possibly stale replica.

  4. Fails

    Evicting a remote marker is as harmless as evicting any other cached key.

    A marker is information, not a copy: its presence says 'the local replica may be stale for this key'. Evicting it means a read may go to the stale replica. Facebook pointed this out explicitly; it is rare enough in practice to accept.

The reasoning

  1. Caches inherit replication lag: invalidate after the data arrives, from each region's own log.
  2. Remote markers send a writer's reads to the master while their write is in flight.
  3. A key whose presence carries meaning isn't a disposable cache entry any more.

Across regions, the cache inherits the database's Replication lag. Two rules follow: invalidate after the data arrives (by tailing each region's own log), and route a writer's reads to the master while their write is in flight.

The marker claim is the subtle one. Stage 1 said any key can be safely evicted. That is true for copies of data. A key whose presence carries meaning is not a cache entry anymore, and the system has to know the difference.

Stage 8 of 9 · Break it

A cluster with an empty cache

Other clusters in the region have warm caches with roughly the same data.

What you need to know first

A cluster serves 2,000,000 reads a second. Warm, its hit rate is 99%. Cold, it's about 5%. About how many database reads a second does it send when cold?

About 1.9 million per second.

Warm: 1% of 2,000,000 = 20,000. Cold: 95% = 1,900,000 a second, 95× what the database is provisioned for.

Other clusters in the region hold roughly the same data, warm. A cold cluster can treat a neighbour's cache as its fallback: on a local miss, fetch from the warm cluster's cache and add it locally, so misses come from memory instead of the database.

A key is written and deleted in both clusters, but the cold cluster receives its delete slightly before the warm one. In that gap, a cold-cluster miss fetches the key from the warm cluster. What happens?

It gets the old value, because the warm cluster hasn't deleted it yet, and adds it to the cold cluster, after that cluster's delete. The stale value now sits in the cold cluster with nothing to remove it.

A short hold-off after each delete (Facebook used two seconds) rejects adds to that key, and a rejected add means "go to the database".

What the stage asks

How do you bring the cold cluster back?

  1. Flawed

    Send it traffic and let the caches fill from the database

    Nearly every request misses, so the cluster's entire read load hits MySQL. With a large cluster, that is a database outage. Facebook said warming this way took days.

  2. Sound

    On a miss, fetch from a warm cluster's cache and add the value locally; deletes in the cold cluster carry a short hold-off that rejects adds

    Misses are served from memory in a neighbouring cluster instead of the database, and the cluster warms in hours. The hold-off (two seconds at Facebook) closes a race: without it, a value deleted after a write could be re-added from the warm cluster before that cluster received the same delete.

  3. Defensible

    Copy a full snapshot of a warm cluster's cache before sending traffic

    It warms everything, including keys that will never be read here, and the data changes while it copies, so you still need invalidation during and after the copy.

What a strong answer covers

  • An empty cache sends its whole read load to the database.
  • A warm peer cache can absorb the misses instead.
  • A value can be deleted locally and then re-added from the peer before the peer gets the delete.
  • A hold-off after deletes rejects adds for a short window, and a failed add means 'go to the database'.

The reasoning

  1. An empty cache sends nearly its whole read load to the database.
  2. Warm a cold cluster from a neighbour's cache instead of the database.
  3. Close the warmup race with a short hold-off that rejects adds after a delete.

Every mechanism in this investigation answers the same question: where does the miss load go? Leases send it to one refiller, gutter sends it to idle servers, warmup sends it to a neighbour's memory. The database only ever sees what it was provisioned for.

Each one also opens a small consistency window, and each closes it with a bounded mechanism (a token, a short expiry, a hold-off) instead of perfect consistency.

Stage 9 of 9 · Defend it

Defend 'best-effort eventual consistency'

Your interviewer: "Your cache can serve stale data in at least four ways. Why not a strongly consistent cache, or write-through updates, and be done with it?"

What you need to know first

"Eventually consistent" is too vague to defend. A strong answer lists:

  • Specific guarantees: for example, a writer sees their own write; a stale value can't be cached after a newer write.
  • Specific windows of staleness: where they come from and what bounds each one.
  • What stronger consistency would cost: usually coordination on the hot path.
  • Where you'd choose differently: data that can't tolerate any staleness.

What would strong consistency for every read cost a cache serving billions of reads a second across regions?

Coordination on the hot path: checking with the source or locking every copy, with cross-region round trips and lost availability during failures

The cache exists to avoid going to the source. Making every read confirm with it removes most of the benefit.

What the stage asks

Defend the consistency model: what is guaranteed, what is not, and what strong consistency would cost here.

Reference answer

What is guaranteed. A user sees their own writes (local delete plus remote markers). A stale value cannot be cached after a newer write's delete (leases). Invalidations are recorded with every commit and replayed if lost, so staleness is bounded, not indefinite.

What is not. Other users may see a slightly old value for a short time: during replication lag, from gutter, or during warmup.

What strong consistency would cost. Every read would have to confirm with the source (or every write would have to synchronously update or lock every copy across clusters and regions) before answering. At billions of reads a second, that means cross-region round trips on the hot path and a cache that stops serving when any participant is unreachable. The whole point of the cache is to avoid the database; strongly consistent caching brings it back into every read.

Why not write-through. Concurrent sets still race across clusters, so ordering is still a problem, and it fills every cluster with values they may never read.

Where I would choose differently. For data where staleness is unacceptable (balances, permissions changes that must apply instantly), read from the master or skip the cache.

What a strong answer covers

  • States what is guaranteed: read-your-writes for the writer, bounded staleness, no indefinite stale sets.
  • Explains what strong consistency costs at this scale: coordination on every read or write, latency, availability during failures.
  • Explains why write-through does not solve ordering and fills caches with unread data.
  • Notes that data needing strong consistency can bypass the cache or read from the master.Supporting

The reasoning

  1. Defend a consistency model with specific guarantees and specific, bounded staleness windows.
  2. Strong consistency at this scale means coordination on every read, which defeats the cache.
  3. Route data that can't tolerate staleness to the master or around the cache.

"Eventually consistent" is not an answer by itself. A strong defence lists the specific guarantees, the specific windows of staleness and the mechanism that bounds each one. That is what Facebook's paper does, and what an interviewer wants to hear.

How you did

Now try it as an interview question

  • “Design a distributed cache.”
  • “How do you keep a cache consistent with the database?”
  • “A hot key expires and the database falls over. What happened and how do you prevent it?”
  • “Design Memcached or Redis as a service for a large company.”

The interview mode mixes stages from this and other investigations with concept recall and questions about your own projects.

Back to the last stage