Cache Invalidation Across a Fleet — a worked solution
Write-through vs write-behind. Two generals.
Try it yourself first.
You will remember almost none of this if you read it cold. The workspace walks the same 10 stages and runs the architecture you draw through a simulator, so you find out where your version breaks before you see ours.
Open the Cache Invalidation Across a Fleet workspaceThe problem
A cache only earns its keep if readers trust it. The instant a write lands in the database, every cached copy of the affected key — across hundreds of cache hosts, in several regions, plus the CDN at the edge — is wrong. Cache invalidation across a fleet is the problem of propagating "this key changed" to all of those copies, quickly, reliably, and provably enough that users do not see stale data.
The hard part is not the happy path. It is that you can never be certain an invalidation was delivered (the Two Generals problem), so a system that assumes "I sent the delete, therefore the cache is correct" is wrong by construction. Production systems instead make invalidation idempotent, replayable from a durable log, version-guarded, and continuously measured — and accept that "correct" means "inconsistent for less than X milliseconds, less than one time in ten billion," not "never stale." This canonical models the look-aside fleet that Meta (memcache + mcsqueal + leases + Polaris), Netflix (EVCache), and Uber (CacheFront + Flux) actually run.
The reference architecture
What each component is for
- App / Client
Issues reads (the overwhelming majority) and the occasional write. A read is a GET the system tries to satisfy from the nearest cache; a write is a mutation that must become durable AND invalidate every cached copy of the affected keys across the fleet.
Why it exists. It is the traffic source and the only thing that actually observes staleness. We model it explicitly because the ~500:1 read:write ratio and the read-your-writes expectation after a mutation are stated from the client's point of view, not the server's.
When it fails. A client that retries aggressively on a slow write turns one mutation into a storm of writes-and-invalidations; the Idempotency-Key plus the gateway rate limit are what keep that from melting the invalidation pipeline.
- Edge CDNFastly (surrogate keys)
Caches GET responses at the edge keyed by URL and tagged with surrogate keys (cache-tags). On a write, the invalidation pipeline issues a surrogate-key purge so every POP drops the affected objects at once.
Why it exists. It is the first 'fleet' that must be invalidated and the cheapest place to absorb read load — a viral key served from the edge never reaches the origin cache or the DB. Including it makes the point that 'invalidation across a fleet' is recursive: the CDN is a second fleet with the same Two-Generals problem, solved the same way (idempotent, replayable, version-tagged purges).
When it fails. A botched purge leaves stale pages at the edge until TTL; a mass re-tag can saturate the purge plane. Detection: edge hit-ratio and purge-latency SLO. Mitigation: surrogate-key granularity discipline plus a synthetic write-then-read canary that pages if an edge purge misses SLO.
- API GatewayEnvoy + WAF
L7 ingress. Terminates TLS, runs the WAF, enforces per-client token-bucket rate limits, validates the bearer token on writes, and routes by intent and geography: cacheable GETs to the regional read service, mutations to the single master-region write service.
Why it exists. Centralizes auth, rate-limit, and geo-routing so a runaway client is shed at the edge before it can storm the write path and the invalidation bus. Folding this into each service would let a crash bypass the limiter and let a retry storm reach the binlog tailer.
When it fails. Gateway saturation reads to users as 'everything is slow'. Detection: active streams, downstream 5xx rate, rate-limit-reject rate. Mitigation: a per-client circuit breaker the on-call can trip, drain mode, per-zone failover.
- Read Service (US)Go look-aside + single-flight
Serves a GET by asking the regional cache fleet for the key. On a hit it returns in ~1ms. On a miss it takes a lease token from the cache, reads the row from the DB, and SETs the value back guarded by that lease — so a concurrent invalidation can reject a stale fill.
Why it exists. A dedicated cache-access tier (vs every app server talking to memcached directly) is where request coalescing / single-flight lives: all concurrent misses for one hot key collapse to a single DB read — exactly what Discord's data services, Uber's CacheFront, and DoorDash's caching library do. It also owns the lease protocol so the stale-set race is fixed in ONE place.
When it fails. If a reader replica's single-flight map is reset (e.g. on deploy), a hot-key miss can stampede the DB. Detection: DB QPS spike correlated with cache miss-ratio. Mitigation: leases at the cache make the herd safe even when a single replica's coalescing fails.
- Write ServiceGo + idempotency
Handles a mutation: writes the new row to the DB of record (the binlog entry is the durable record of the change), then issues an immediate best-effort invalidation to the bus so the local fleet drops the key fast. It deliberately does NOT trust that best-effort message to arrive.
Why it exists. Splitting writes from reads keeps the read tier stateless and lets writes carry the heavier machinery: idempotency, master-region routing, and dual-path invalidation. The best-effort publish is the low-latency leg ('write-through-ish'); correctness comes from the binlog tailer, not from this publish.
When it fails. If the writer crashes after the DB commit but before the best-effort publish, the local cache stays stale until the binlog tailer catches it (tens of ms) or the TTL expires (worst case). This is the Two-Generals gap — bounded by the reliable path plus TTL, never zero.
- Cache Fleet (US)Memcached + mcrouter
Hundreds of memcached hosts behind mcrouter, consistent-hash-sharded by key. Serves gets in ~1ms, hands out lease tokens on misses, applies deletes (invalidations) to the owning host, and stamps every entry with a version so an older value can never overwrite a newer one.
Why it exists. This IS the fleet the problem is named for. It exists to shield the DB from ~98% of reads; everything else in the diagram exists to keep these hosts from serving wrong data. mcrouter is what makes invalidation TARGETED — a delete routes to the one host that owns the key instead of broadcasting to all ~390 hosts in the region (which would be ~390x the traffic).
When it fails. A dead host's keys miss through to the DB until the in-fleet gutter hosts or a rehash cover them — bounded blast radius (1/N of the keyspace), no data loss. The dangerous failure is SILENT staleness: a lost invalidation leaves a host serving a stale value until TTL, invisible from the read path — caught only by the consistency monitor and the independent canary. RTO ~5s for a host; RPO 0 (no durable data here).
- DB of RecordMySQL (sharded, binlog)
The durable, sharded relational store. Every committed mutation lands here first and is written to the per-shard binary log; that binlog is BOTH the replication stream and the source of invalidations — the cache tailer derives 'drop key X' from 'row X changed'.
Why it exists. Something has to be the source of truth, and it has to produce an ordered, durable, replayable log of changes. Using the binlog as the invalidation source — rather than trusting the application to publish — is the single most important correctness decision: the log cannot 'forget' a committed write, so invalidations can always be replayed (Meta's mcsqueal, Uber's Flux, and Debezium/Maxwell/Canal all do exactly this).
When it fails. Primary loss blocks writes for the failover window (RTO ~30s via Orchestrator/Patroni); reads continue from replicas (eventual). If the binlog tailer cannot keep up, invalidations queue behind commits and fleet-wide staleness grows with the lag — the most insidious failure, because the read path still looks healthy. Detection: binlog-tail lag vs commit rate (page at >30s lag, hard-page at >6h, well before the 72h retention cliff).
- Invalidation Pipelinemcsqueal-style binlog tailer
Reads the DB's per-shard binlog, extracts the cache keys affected by each committed change, batches them, and publishes durable invalidation messages to the bus. Because it derives invalidations from the committed log, it cannot miss a write the application forgot to invalidate.
Why it exists. This is the answer to Two Generals. You can never confirm an invalidation was delivered, so you do not rely on best-effort delivery — you derive invalidations from a durable, ordered log that can be replayed from any offset. It exists precisely so the writer's best-effort publish is ALLOWED to be lossy.
When it fails. If a tailer falls behind binlog retention it cold-starts and must re-scan — during which the fleet is broadly stale. Detection: per-shard tail lag with a concrete SLO — page at >30s lag, hard-page at >6h, both well before the 72h retention cliff. Mitigation: parallelize per shard; coalesce harder under load; backpressure writes only as a last resort.
- Invalidation BusKafka / Wormhole
The pub/sub backbone that carries invalidation messages from the writer (best-effort) and the binlog tailer (reliable) to every cache fleet, in every region, plus the consistency monitor. Partitioned by key-hash so all invalidations for one key stay ordered.
Why it exists. Fan-out and decoupling: one invalidation must reach the local fleet, every remote-region fleet, and the monitor, without the producer knowing who the consumers are or blocking on them. A durable, replayable log (vs a fire-and-forget multicast) is what lets a consumer that was down replay the invalidations it missed on recovery.
When it fails. A bulk update (migration, backfill) can still enqueue millions of per-key deletes and blow out consumer lag, widening staleness globally — which is why the upstream pipeline coalesces by key before publishing here. Detection: per-partition consumer lag plus producer/consumer rate gap. Cert expiry on the bus silently stops ALL invalidations — alert on cert-expiry-days independent of the pipeline. Leader election on a partition is ~10s (RF=3 ISR), during which that partition's deletes queue but are not lost.
- Consistency MonitorPolaris-style
Treats the cache as a black box and continuously checks the invariant 'the cache eventually agrees with the DB'. It subscribes to the invalidation stream, samples the committed version in the DB, compares it against what the fleet is serving, and re-publishes a corrective invalidation for any entry still stale past the window.
Why it exists. Because Two Generals guarantees the pipeline is never perfect, the only durable defense is to MEASURE inconsistency and treat it as an SLO — not to assume the pipeline is correct. Meta's Polaris did exactly this and took cache consistency from six nines to ten nines (<1 stale entry per 10 billion within 5 minutes). It is also the backstop that drains and replays the dead-letter queue once a poison message is fixed.
When it fails. If the monitor itself is down you go BLIND — staleness can climb with no alert (the dead-man's-switch problem), and because it both consumes and re-publishes on the bus, a bus outage can silence it and the pipeline together while the staleness SLI flatlines deceptively green. Detection therefore CANNOT live here: the independent Canary node is the real dead-man's-switch, plus a 'monitor-not-reporting' alert that fires if this service stops emitting the staleness SLI for >60s.
- Dead-Letter QueueKafka DLQ topic
Holds invalidation messages that could not be parsed or applied — malformed records, schema mismatches, repeated apply failures — so the main pipeline keeps flowing and the consistency monitor (or a human) can triage and replay them.
Why it exists. Head-of-line protection. Without a dead-letter path, one poison message wedges its partition and every key behind it goes stale. The DLQ converts 'one bad message stalls a region' into 'one bad message is parked and alerted' — the difference between a page and an outage.
When it fails. A growing DLQ means invalidations are silently NOT being applied — stale data is accumulating. Detection: DLQ depth plus age-of-oldest-message alarms. A DLQ filling faster than it drains is an incident, not a backlog.
- DB Replica (EU)MySQL async replica
A read-only replica of the DB of record in a second region. It applies the master's binlog asynchronously and serves local reads on cache misses, so an EU read does not cross the Atlantic to the US master.
Why it exists. Latency and survival. Without a regional replica, every EU cache miss would pay a cross-region round trip, and a US-region outage would take everyone down. It is also the DR target — promotable to master if the primary region is lost.
When it fails. Replication lag here is directly visible to EU users as stale reads after a US write ('I changed it but EU still shows old'). Detection: cross-region replication-lag metric plus a per-region staleness SLO. Mitigation: the remote marker forwards read-your-writes reads to the master until the replica catches up.
- Cache Fleet (EU)Memcached + mcrouter
A full, independent memcache fleet in the second region serving EU reads at ~1ms. It receives invalidations from the bus (cross-region) and fills misses from the EU DB replica.
Why it exists. Each region needs its own fleet — caching is pointless if every read crosses an ocean. Modeling it explicitly surfaces the hardest correctness bug in the whole design: a cross-region invalidation can RACE the cross-region data replication, so the fleet must not refill from a replica that has not caught up yet.
When it fails. If the version guard is missing, the classic bug appears: invalidation arrives before the replicated write, the cache evicts, a read refills from the still-stale replica, and the fleet re-caches stale data no further invalidation will clear. Detection: Polaris flags it as PERMANENT (not transient) staleness. RTO ~5s/host, RPO 0.
- Read Service (EU)Go look-aside + single-flight
The EU region's cache-access tier. Same look-aside + single-flight + lease logic as the US reader, but reads fill from the EU DB replica — except when a key's remote marker is set, in which case it forwards the read to the US master for read-your-writes.
Why it exists. A region needs a local read tier for the same reasons the primary does (coalescing, leases, negative caching). It is modeled separately to make read-your-writes-across-regions concrete: the remote marker lives here, on the read path, where the decision 'is my local replica fresh enough?' is actually made.
When it fails. If the remote marker is wrong or missing, EU users see stale data right after their own writes (read-your-writes violation) — the most common user-visible cache bug in a multi-region system. Detection: a per-region read-your-writes canary. Mitigation: a conservative marker TTL tied to the p99 replication lag.
- Synthetic Canaryout-of-band prober
Continuously writes a probe mutation through the real front door (gateway -> write service -> DB -> invalidation) and then reads the probe key back directly from every cache fleet, in every region. If a probe write is not reflected fleet-wide within the staleness SLO, it pages — regardless of what any internal metric says.
Why it exists. The system's single worst failure is SILENT: invalidations stop, reads still return 200s, and staleness climbs with zero 5xx. The consistency monitor that measures this consumes and re-publishes on the SAME bus it is watching, so a bus-wide failure can silence the monitor AND the pipeline together while the dashboard stays green. A watchdog cannot share a failure domain with the thing it watches — so the canary lives OUTSIDE the bus and exercises the end-to-end path itself.
When it fails. If the canary itself is down you lose the independent signal — so it runs 3 replicas across zones and its own liveness is alerted by an even-simpler external uptime check (the watchdog's watchdog). It is intentionally the dumbest, most independent component in the system precisely so it almost never fails for the same reason everything else does.
Stage by stage
The same 10 stages the workspace walks, answered.
01Clarifications
What would you ask before drawing a single box?
Questions worth asking before drawing a box:
- What is the source of truth? Here: a sharded relational DB. The cache is disposable; losing a cache host loses nothing durable.
- Look-aside or read-through? Look-aside (the app/read-tier reads cache, falls through to the DB on miss, and writes back). This is where the stale-set race lives, so it drives the lease design.
- Invalidate or update on write? Invalidate-by-delete is the safe default (a delete is reorder-safe). Write-through (update the cached value) and write-behind (buffer in cache, flush async) are the alternatives we discuss in deep-dives.
- How fresh must reads be? We target a bounded staleness window, not linearizability. Read-your-own-writes is required; other users tolerate sub-second staleness.
- Single-region or multi-region? Multi-region with a single master region for writes (single-writer-per-shard) and read replicas + a full cache fleet per region.
- Read:write ratio? ~500:1 (TAO's published figure). This is why almost all design pressure is on keeping reads correct, not on write throughput.
Assumptions:
- One service tier doing ~10M cache gets/sec, ~20k writes/sec, 98% hit ratio, 5 regions, ~60 TB hot working set.
- Staleness tolerated: ≤ ~160 ms same-region, ≤ ~500 ms cross-region. TTL backstop = 300 s caps staleness from lost invalidations.
02Functional reqs
What must this system actually do?
- Read a key: serve from the nearest cache; on miss, fill from the DB and populate the cache.
- Write a key: durably persist the new value AND invalidate every cached copy of the affected keys, fleet-wide and cross-region.
- Invalidate-by-delete on write (not value-replace), so the next read re-derives the fresh value.
- Read-your-own-writes: a client that just wrote must not subsequently read its own stale value, even cross-region.
- Replay invalidations: a cache fleet that was unreachable must catch up from the durable invalidation log on recovery.
- Measure consistency: continuously report the stale-read ratio as an SLI, and auto-correct entries found stale past the window.
- Bound worst-case staleness with a TTL even when every invalidation for a key is lost.
03Non-functional
What must it promise about speed, uptime and correctness?
SLOs (one service tier):
| Dimension | Target |
|---|---|
| Read latency (cache hit) | p50 1 ms, p99 5 ms, p99.9 20 ms |
| Read latency (miss → DB fill) | p99 ~12 ms; blended p99 ~5.1 ms at 98% hit |
| Invalidation delivery (commit → all caches, same region) | p50 8 ms, p99 45 ms, p99.9 160 ms |
| Invalidation delivery (+ cross-region) | p50 +70 ms, p99 +150 ms, p99.9 +400 ms |
| Staleness window SLO | ≤ 160 ms same-region, ≤ 500 ms cross-region |
| Consistency (steady state) | ≤ 1 inconsistent entry per 10^10 within 5 min (Polaris) |
| Read availability | 99.99% (reads are the product) |
| Invalidation-pipeline durability | invalidations are acks=all, RF=3 — we refuse to lose one |
Error budget. Read availability 99.99% → ~52 min/yr. The dominant user-visible failure is not unavailability but silent staleness, so the budget that actually governs on-call is the staleness SLO burn rate: page on a fast burn (2% of the staleness budget in 1 h, confirmed by a short window), ticket on a slow burn (SRE Workbook multi-window multi-burn-rate alerting). A gray failure — invalidations silently stopping while reads look healthy — must be caught by the staleness SLI, not by a 5xx alarm.
04Capacity estimation
How much load and data does this have to hold?
Start from the anchored load and derive every shard/replica/fleet number rather than guessing.
1. Throughput. Anchor reads at R = 10,000,000 gets/sec for one service tier (Meta's memcache served billions/sec fleet-wide; this is one tier of that). With TAO's 500:1 read:write, writes/invalidations W = R/500 = 20,000/sec. At a 98% hit ratio, DB-bound misses = R·(1−0.98) = 200,000/sec.
2. Fleet sizing. Hot working set D = 60 TB. Hosts are 256 GB RAM, ~192 GB usable after overhead. With 25% headroom: hosts/region = 60 TB · 1.25 / 0.192 TB ≈ 390. Across 5 regions → ~1,950 cache hosts. Per-host get load = 10M / 390 ≈ 25,600 gets/sec/host — comfortably under the ~1M ops/sec single-host ceiling.
3. Why targeted, not broadcast. A write must reach the hosts holding the key. Targeted (mcrouter consistent-hash → only the owning host, in each region): fan-out ≈ regions = 5, so W·5 = 100,000 invalidations/sec. Broadcast (every host, e.g. for a wildcard/surrogate purge): fan-out = 1,950, so W·1,950 = 39,000,000 deliveries/sec — 390× more. This single multiplier is why the pipeline tails the binlog and routes through mcrouter to the owning host instead of multicasting to the fleet.
4. The TTL backstop is cheap. A 300 s TTL caps how long a lost invalidation can serve stale data; its cost is the extra refills it forces. The naive bound — every one of 60 TB / 1 KB ≈ 6×10¹⁰ entries expiring and refetching once per window — would be 6×10¹⁰ / 300 s ≈ 200,000,000 refills/sec, which is absurd (20× the entire 10M/sec read rate). It never materializes, because a TTL only triggers a refill when a still-wanted key is requested again after it expires, and those re-requests are already part of the organic miss stream (R·(1−0.98) = 200,000/sec); the cold majority of the 60 TB is never re-requested inside any 300 s window. So the backstop rides on the existing ~200,000/sec miss rate for essentially nothing while capping lost-invalidation staleness at 5 minutes. Drop TTL to 60 s and a continuously-hot key refills ~5× more often (every 60 s vs 300 s) and it starts to bite; raise it to an hour and a lost invalidation hides for an hour.
5. Where the architecture flips.
| Decision | Flips when |
|---|---|
| Broadcast → targeted (mcrouter) | fleet × W exceeds router capacity — here at ~50+ hosts |
| App-publish → binlog-tail (mcsqueal) | multiple writers / replication make app-publish lossy; correctness needs the DB log as source of truth at any real scale |
| No-lease → leases | any hot key with > ~1k concurrent miss-refills (herd cost > lease cost) |
| Single Kafka topic → partitioned | one topic saturates ~100k–1M msg/s; at W·5 = 100k/s you partition by key-hash before the cap |
| TTL-only → explicit invalidation | short-TTL refills overload the DB; long TTL + invalidation wins above ~1M gets/sec |
| Write-back → write-through / invalidate | staleness SLO is tighter than the async flush window, or the data is durability-sensitive |
05API design
What does the outside world call, and what comes back?
Read — two classes, by data sensitivity:
Public / non-PII object (CDN-cacheable, no auth):
GET /v1/objects/:key
200 OK
Cache-Control: public, max-age=30, stale-while-revalidate=30
Surrogate-Key: obj:123 owner:42
ETag: "v=8814" # version/generation, used by the cache version-guard
{ "key": "obj:123", "value": { ... }, "version": 8814 }
Private / tenant / PII object (auth required, bypasses the shared CDN):
GET /v1/objects/:key
Authorization: Bearer <token> # or session; tenant scope authorized at the gateway/read service
200 OK
Cache-Control: private, no-store # never lands in the shared Fastly edge cache
ETag: "v=8814"
{ "key": "obj:123", "value": { ... }, "version": 8814 }
Tenant scope is enforced by the gateway/read service from the token, not by key-secrecy — knowing :key is not authorization. The obj:123 example above is a public, non-PII object: only that class may be served Cache-Control: public through the third-party edge. PII-bearing objects use the private form (Cache-Control: private / no-store) and bypass the CDN entirely, so they are never co-resident in a shared public edge cache. Internal cache-fleet entries are still version-guarded identically regardless of class.
Write (mutation, master region):
POST /v1/objects/:key
Authorization: Bearer <token>
Idempotency-Key: 7f1c... # dedups retried mutations + invalidations
{ "value": { ... } }
200 OK
{ "key": "obj:123", "version": 8815 }
# Side effects, in order:
# 1. row committed to DB (binlog entry @v8815) — durable
# 2. best-effort DELETE published to the bus — fast, lossy
# 3. binlog tailer derives DELETE @v8815 — reliable, replayable
Invalidation message (on the bus):
{ "key": "obj:123", "op": "delete", "version": 8815, "region": "us-east", "source": "binlog|app|reconciler" }
A cache applies a delete only if version ≥ stored.version (older never overwrites newer); a lease-gated SET is rejected if its lease was invalidated. CDN purge is the same operation by Surrogate-Key.
Lease protocol (look-aside miss): GET key miss → cache returns (miss, lease_token); reader reads DB, then SET key value WITH lease_token; the SET is dropped if a DELETE invalidated lease_token in between. A token is re-issued at most once per ~10 s per key, which throttles the herd.
06Data model
What gets stored, and what is it looked up by?
Cache entry (memcached value): value, version (monotonic per key), lease_token?, negative? flag, TTL. The version is the load-bearing field — it makes out-of-order and duplicate invalidations safe and makes staleness detectable.
DB of record (sharded by entity_id): the canonical rows. Each shard's binlog is the ordered, durable change log — the contract the whole invalidation pipeline is built on. Retained 72 h so a tailer can fall behind and replay.
Remote marker (per key, in the writing user's home-region cache pool, not the master): a short-lived flag set there when that region's write is forwarded to the master, so the read service's per-miss marker check is a local ~1 ms lookup and only reads of marked keys pay the cross-region hop to the master; it clears once replication catches up (conservative TTL tied to p99 replication lag), restoring read-your-writes.
Database choice — recommended. A sharded relational store (MySQL with binlog) for the record, memcached + mcrouter for the fleet, Kafka/Wormhole for the invalidation bus. Why MySQL: invalidation correctness rides on an ordered, durable, replayable per-shard log, and the binlog is exactly that. ScyllaDB/Cassandra (Discord) or DynamoDB are fine record stores if you tail their change streams (CDC / DynamoDB Streams) instead of a binlog. Redis/Valkey replaces memcached when you want richer data structures or built-in cross-region replication (EVCache). The one non-negotiable: the source of invalidations must be the DB's committed change log, not the application's best intentions.
07High-level design
Which components handle a request, and in what order?
Read path (the product, ~98% of traffic). Client → CDN → (edge miss) → Gateway → Read Service → Cache Fleet. On a cache hit the Read Service returns in ~1 ms. On a miss it takes a lease, reads the DB, and SETs the value back under the lease — single-flight collapses a herd of concurrent misses into one DB read. This e1 Client→CDN read is unauthenticated and only carries the public/non-PII read class; private/PII reads are authenticated (Cache-Control: private/no-store) and must skip the shared edge — that needs a direct authenticated Client → Gateway read edge (not modeled in the current graph, which only routes private reads through the CDN-fronted e1/e2).
Write path (the hard part). Client → Gateway → Write Service → DB of Record — writes take the direct, bearer-authenticated Client → Gateway edge and bypass the edge cache entirely (the CDN caches GET/HEAD only). The commit is the durable truth. The Write Service then fires a best-effort invalidation to the bus (low latency, allowed to be lost). Independently, the Invalidation Pipeline tails the DB binlog and publishes the reliable, versioned invalidation. Two paths, on purpose: the fast one for latency, the durable one for correctness. This is "write-through vs write-behind" made concrete — the best-effort publish behaves like an eager write-through of a delete; the binlog tail behaves like a guaranteed write-behind that can never forget.
Fan-out. The Invalidation Bus delivers each delete, partitioned by key-hash (per-key order), to the local Cache Fleet (targeted via mcrouter), every remote-region fleet, the Edge CDN (as a surrogate-key purge), and the Consistency Monitor. A dead-letter queue parks poison messages so one bad record cannot wedge a partition.
Cross-region. The DB binlog replicates to the EU DB Replica (async); the bus delivers invalidations to the EU Cache Fleet; the EU Read Service serves local reads but honors the remote marker to forward read-your-writes reads to the master. The EU fleet's version guard rejects a refill from a still-stale replica — defeating the invalidation-beats-replication race.
Reconciliation. The Consistency Monitor (Polaris-style) watches the invalidation stream, samples committed versions in the DB, and re-publishes a corrective invalidation for anything still stale — turning Two-Generals uncertainty into a measured, bounded SLO and draining the DLQ on recovery. Because it shares the bus's failure domain, the Synthetic Canary sits outside the pipeline entirely: it writes a probe through the real front door and reads it back directly from each fleet, paging if a probe is not reflected fleet-wide within SLO — the independent dead-man's-switch that catches the silent-staleness gray failure the monitor itself could be blind to.
08Deep dives
Which part breaks first, and what do you do about it?
Two Generals & the dual path
You can never confirm an invalidation was delivered, so don't build on confirmation. The best-effort publish from the Write Service is fast and lossy; the binlog tailer derives the same invalidation from the durable, ordered commit log and can replay it from any offset. Note this is log-tailing / CDC (the binlog is the ordered event log — an event-sourced-style projection feeds the caches), not a transactional-outbox table written inside the business transaction; we get at-least-once delivery deduped by the version guard, not exactly-once from an outbox row. Delivery is made idempotent (delete is reorder-safe) and version-guarded (older never overwrites newer). The TTL caps worst-case staleness from total loss. The Consistency Monitor measures what leaks through — but it shares the bus's failure domain (it consumes and re-publishes on the same bus), so the Synthetic Canary, living entirely outside the bus, is the actual dead-man's-switch. No single mechanism is sufficient; together they bound staleness to "<1 in 10^10 within 5 min" rather than promising zero.
Stale-set race & leases
The canonical look-aside bug: reader A misses, reads v1, stalls; writer commits v2 and invalidates; slow A now SETs v1, poisoning the cache until the next write or TTL. The fix is leases (NSDI'13): on a miss the cache hands A a token; a concurrent delete invalidates outstanding tokens, so A's late SET is rejected. The same token, re-issued at most once per ~10 s per key, throttles the thundering herd — one refiller, everyone else briefly serves stale or waits.
Write-through vs write-behind vs invalidate
Invalidate-and-refill (default): write DB, delete cache, next read refills. Cheap, reorder-safe, costs one miss. Write-through: write DB and SET the new value into the cache synchronously — gives instant read-your-writes but pays write latency and reintroduces the stale-set race (must be version-guarded), and a value-SET can lose a concurrent newer write. Write-behind (write-back): acknowledge into the cache and flush to the DB asynchronously — lowest write latency, but a crash loses acknowledged writes (unacceptable for durability-sensitive data) and the cache becomes a source of truth you must replicate. We invalidate by default and reserve write-through for keys that demand zero-miss read-your-writes.
Targeted vs broadcast fan-out
Broadcasting every invalidation to all ~1,950 hosts is 39M deliveries/sec — 390× the targeted rate. mcrouter's consistent-hash routing sends each delete to the one host that owns the key, per region. Broadcast is reserved for wildcard/surrogate purges where the key→host mapping is unknown (CDN cache-tags), and even there the CDN's purge plane fans out hub-to-spoke, not N×N.
Cross-region read-your-writes & the invalidation race
Two cross-region hazards. (1) Lag: a US write reaches EU seconds later, so an EU read sees stale data — fixed by the remote marker routing RYW reads to the master until replication catches up. (2) The race: a global-broadcast invalidation can arrive in EU before the replicated write, the EU fleet evicts, and a refill re-caches the still-stale replica value — fixed by version-guarded fills (reject a fill whose replica version is older than the invalidation's version). Meta's alternative is to derive each region's invalidations from that region's replicated binlog, so an invalidation can't outrun its data; we model the global-bus + version-guard approach (Netflix EVCache) and treat the per-region-binlog approach as the rejected alternative.
Hot keys & thundering herd
Traffic is Zipfian; one key can take 100k gets/sec. On expiry, 100k concurrent misses would crush the DB shard. Leases collapse that to one DB fetch; single-flight at the Read Service collapses concurrent in-process misses; TTL jitter (±10%) stops synchronized expiry; hot keys can be replicated across multiple cache keys or promoted to a per-region replica. The CDN absorbs the very hottest keys before they reach the origin at all.
Observability
SLIs (RED/USE per component): Read Service — request rate, error rate, p99 (RED); Cache Fleet — hit ratio, get p99, eviction rate, per-host connection saturation (USE); Bus — produce/consume rate, per-partition consumer lag, DLQ depth + age-of-oldest; Pipeline — per-shard binlog-tail lag; DB — commit rate, replication lag (in-region and cross-region). The headline SLI is the stale-read ratio from the Consistency Monitor, reported at 1/5/10-minute windows. But a measured SLI can lie by going silent — if the monitor dies, the stale-read ratio flatlines deceptively green. So two SLIs are independent of the monitor: the Synthetic Canary's probe-not-reflected-within-SLO signal (which exercises the real path from outside the bus) and a "monitor-not-reporting" alert that fires if the Consistency Monitor stops emitting its SLI for >60 s. SLO burn-rate alerts page on fast staleness burn and on invalidation-pipeline lag (binlog-tail >30 s, hard-page >6 h) before it becomes user-visible. The six dashboards on-call opens: (1) staleness SLI by region + timescale — with a freshness check that it's still updating; (2) invalidation pipeline lag (binlog tail → bus → apply); (3) cache hit-ratio & DB miss-QPS (stampede detector); (4) bus consumer lag + DLQ depth/age; (5) cross-region replication lag + remote-marker rate; (6) canary probe latency + CDN purge latency (the two fleet-wide edge signals). Tracing: the invalidation path is traced end-to-end (commit → tail → bus → apply) with head-based sampling (~1%) plus 100% sampling of the canary's write-then-read so a missed invalidation always has a trace.
Deploy & runtime
Rolling, never flush-all: a cache-tier deploy restarts hosts in small batches so the fleet never goes cold at once (a cold flush stampedes the DB). Cold-cluster warmup: a new/cold cluster reads through a warm one for hours before taking full traffic. Schema/format migrations of the invalidation message are backward-compatible (new fields optional) and rolled out consumer-first. Consumer-group rebalance on the bus is incremental (cooperative-sticky) so a deploy doesn't pause invalidation delivery. Secrets: the bus mTLS certs auto-rotate (cert-manager/SPIFFE) with an expiry-burn alarm independent of the pipeline — an expired cert silently stops all invalidations.
Multi-region / DR
Posture: active-active reads, single master region for writes (active-standby for the write path). Each region has a full cache fleet + DB replica and serves local reads. Writes route to the master region (single-writer-per-shard) so there are no conflicting cross-region invalidations. DR: losing a cache fleet is a non-event (reads miss to the local DB replica; RPO 0). Losing the master region is a deliberate promotion of a replica to master. Be honest about the two different RTOs here: the DB-promotion RTO is ~300 s (db_replica.failover.rtoSeconds) — the time to promote a replica and re-point the write service — but the user-facing region-failover RTO is realistically 5–15 min, because DNS TTLs, resolver caching, and sticky clients keep sending writes to the dead region long after promotion completes. Claiming a single 300 s number is how a postmortem ends up explaining why the flip "took 12 minutes" against a 5-minute SLO. RPO ~5 s (cross-region async replication lag — you lose up to 5 s of acknowledged master writes). Promotion is manual so a flapping link can't trigger a region flip. The invalidation bus is regional with cross-region mirroring, so a region loss doesn't lose the durable invalidation log.
Security
AuthN/Z: browser→edge is unauthenticated for cacheable GETs (so the CDN can serve them) and bearer-token on writes; every internal hop is mTLS. Tenant isolation: cache keys and surrogate keys are tenant-scoped so one tenant can't purge or read another's keys. PII: values may contain PII, so cache hosts are in-VPC, encrypted in transit, never persisted to disk (no RDB/AOF), and short-TTL'd. Threat headlines: (1) cache poisoning — an attacker-influenced key or a stale-set writes bad data; mitigated by version-guards + leases + immutable-once-written keys; (2) purge abuse — a forged surrogate-key purge as a DoS; mitigated by authenticating the purge API (e23 is auth: api-key) and rate-limiting it; (3) negative-cache poisoning — caching a 404 for a key that then exists; mitigated by short negative TTL + explicit invalidation of negatives on create. Vendor boundary: the Edge CDN is a managed third party (Fastly/Cloudflare) — its purge API has its own SLA and rate limits, so its surrogate-key purge is a best-effort fleet whose guarantee is bounded by the vendor and backstopped by the short edge TTL, exactly the Two-Generals posture applied one layer out.
09Trade-offs
What did this design cost, and what breaks at 10×?
Trade-offs we accepted (and the alternative we rejected)
- Look-aside + invalidate-by-delete over write-through-by-value. Rejected write-through-by-value because a value-SET races concurrent writes and reintroduces the stale-set bug; delete is reorder-safe. Cost: one refill miss per write.
- Dual-path invalidation (best-effort + binlog tail) over app-publish-only. Rejected app-publish-only because it's lossy and can't replay; the binlog can't forget a committed write. Cost: a second pipeline to run and monitor.
- Global bus + version-guard for cross-region over per-region-binlog-tail. Rejected nothing outright — both are real (EVCache vs Meta); we chose the global bus for operational simplicity and pay for it with the version-guard that defeats the invalidation-beats-replication race.
- Single master region over multi-master. Rejected multi-master because conflicting cross-region writes create conflicting invalidations needing reconciliation (CRDTs/LWW); single-writer-per-shard sidesteps it. Cost: cross-region write latency + a remote-marker for RYW.
- Eventual cache consistency with a measured SLO over strong consistency. Rejected strong (lease-everything / synchronous fan-out) because it would couple read latency to fleet-wide delivery; we accept bounded staleness and measure it. Cost: a Consistency Monitor and the admission that we're never at zero.
- TTL backstop over invalidation-only. Rejected invalidation-only because Two Generals guarantees some loss; TTL caps it. Cost: organic refills at TTL.
Failure modes & mitigations
| Failure | Detection signal | Blast radius | Mitigation |
|---|---|---|---|
Lost invalidation (Two Generals) — e8/e11 | staleness SLI burn | one key/host, until TTL | binlog tailer (cdc) replay + version stamp + TTL + Polaris |
Stale-set race — reader/cache | Polaris flags recently-written keys | per-key, sticky | leases (cache) reject the late SET |
Thundering herd — cache/db | DB QPS spike vs miss-ratio | the DB shard, can cascade | leases + single-flight + TTL jitter |
Invalidation storm — broker | consumer lag, rate gap | every region, global staleness | coalesce by key, rate-limit/prioritize, partition headroom |
Cross-region lag — db_replica | replication-lag + per-region staleness SLI | a geography for the lag window | remote marker (reader_replica) → master |
Cross-region race — cache_replica | Polaris: permanent staleness | per-key in the replica region | version-guarded fills |
Poison message — broker/cdc | consumer crash-loop, lag w/ no progress | the partition, head-of-line | schema-validate + DLQ (dlq) + skip-and-alert |
Tailer behind binlog — cdc | per-shard tail lag (page >30 s, hard >6 h) | fleet-wide staleness ∝ lag | parallelize per shard; coalesce harder; alert well before 72 h retention |
Cold cache on deploy — cache | hit-ratio cliff at rollout | whole tier → DB | rolling restart + warmup + in-fleet gutter hosts |
Cert expiry on bus — broker | cert-expiry-days; throughput → 0 | global silent staleness | auto-rotation + expiry-burn alarm + canary |
Consistency monitor down (dead-man's-switch) — reconciler | canary probe miss + "monitor-not-reporting" >60 s | total silent staleness, no 5xx | independent canary outside the bus + monitor liveness alert |
Stale edge after write — cdn | CDN purge-latency SLO + canary | edge-cached objects until TTL | broker → cdn surrogate-key purge (e23) + short edge TTL |
Open questions for human review
- CDN purge plane is now wired (
broker → cdn,e23) but modeled as a single managedkind: cdnfleet rather thankind: external. We keep the dedicatedcdnkind (consistent with the rest of the catalog and its rich edge-TTL/surrogate-key internals) and capture the third-party/vendor-SLA nature in prose; at higher edge scale it deserves its own coreless-purge relay component (Cloudflare CacheDB). Confirm the kind choice. - Gutter pool is modeled as a reserved sub-pool within the
cachefleet node (idle hosts, ~2% capacity) rather than a separate tier, matching NSDI'13. Confirm that's the right altitude vs a dedicatedgutternode. - Per-region binlog-tail vs global bus for cross-region invalidation is presented as a chosen trade-off; a reviewer who values the "invalidation can't outrun its data" guarantee may prefer the Meta per-region model, which would add a tailer + bus per region.
- Write-through for a hot-RYW subset (e.g. a user's own profile) is discussed but not a separate path; confirm whether a small write-through lane is worth the extra complexity vs the remote marker alone.
- DLQ drain ownership is the Consistency Monitor here; some teams prefer a separate human-gated replay tool to avoid an automated replay re-introducing a fixed-then-unfixed bug.
Primary sources
- Scaling Memcache at Facebook (NSDI 2013)
- TAO: Facebook's Distributed Data Store for the Social Graph (ATC 2013)
- Cache Made Consistent / Polaris (Meta, 2022)
- Netflix EVCache global replication
- Uber CacheFront
Now defend it
Reading a design is not the same as being able to hold one under questioning. The workspace asks the same questions an interviewer would, and the simulator disagrees with you when the diagram does not support the claim.
Work Cache Invalidation Across a Fleet yourselfMore in Caching, Proxies & the Edge
Everything between the client and the origin: in-memory caches, CDNs, load balancers and service proxies — and the three ways a cache betrays you.
- Build Build RedisAn in-memory data-structure server: one thread, rich types, optional persistence, async replication. Internalize the cost of single-threaded simplicity and a dozen caching/HA decisions get easier.
- Build Build a CDNA globally-distributed reverse proxy whose only job is to (a) terminate the user's TCP/TLS milliseconds away and (b) serve a cached origin response so origin never sees the request. Internalize edge caching, anycast, TTL, revalidation, SWR, purge, the Vary footgun, origin shield, bypass, and hit ratio — and the dozen ways to misconfigure each.
- Build Build a Service Mesh (Envoy / Istio style)Every microservice request crosses two proxies. This curriculum is what they do: routing, load balancing, timeout-and-retry-budget, circuit breakers, outlier detection, token-bucket rate limits, mTLS with workload identity, and a control plane that streams config to all of them. Build it in the order the production problems show up — and feel why Envoy plus a control plane has eaten the east-west world.
- Build Build a gRPC-style RPC frameworkEvery microservice talks over RPC, and the framework you ship determines half the system's failure modes. Build an RPC framework with codec, streams, deadlines, cancellation, retries, interceptors, and load-aware client-side balancing — and feel why gRPC ate the polyglot RPC market and why Thrift and JSON-over-HTTP linger.