Technologies that implement this pattern: Redis · PostgreSQL · Elasticsearch · Cassandra · DynamoDB · Apache Kafka
Why This Matters#
Read scaling is where most production systems quietly spend their money. The typical consumer product runs somewhere between 10:1 and 1,000:1 reads to writes. A timeline is written once and read thousands of times. A product page is edited weekly and viewed millions of times a day. If you get writes wrong, you lose data. If you get reads wrong, you lose the budget, the latency SLO, and eventually the primary database — at 9am on the busiest day of the year.
Most candidates treat read scaling as a caching question: "put Redis in front of it." Staff engineers treat it as a staleness-budget question. Every read-scaling technique — CDN, cache, replica, materialized view, search index, precomputed feed — is a copy of the truth that is some number of milliseconds or seconds behind. The design question is not which cache. It is how stale can each read be, who signed off on that, and what happens to the copy when the truth changes.
The second reframe: read scaling is a tail problem, not an average problem. A 99% hit ratio sounds great until you notice the 1% of misses all land on the same celebrity profile, at the same second, after the same TTL expiry. Hot keys, stampedes and cold starts — not average QPS — decide whether the primary survives.
If you can walk an interviewer from "what's the staleness budget per read path" to "which copy serves it" to "how the copy is invalidated and who gets paged when it's wrong," you are answering at Staff level.
The 60-Second Version#
- Classify every read path by staleness budget first. Most reads tolerate 1–60 seconds of staleness; a few (balances, inventory at checkout, permissions after revoke) tolerate ~0. Those two classes get different architectures — never one blanket cache.
- Hit ratio is exponential leverage. Going from 95% → 99% hit ratio cuts origin load 5× (5% → 1% of traffic). Going from 99% → 99.9% cuts it another 10×. Measure
cache.hit_ratioper key-prefix, not globally. - Read replicas buy ~N× throughput but add lag. Postgres/MySQL async replicas typically lag 10–500ms, and seconds under write bursts or long queries. Any read-after-write path needs routing to the primary or a lag-aware session token.
- Cache-aside with TTL + jitter is the default. TTL ±10–20% jitter prevents synchronized expiry. Pair it with request coalescing (singleflight) so one miss produces one origin query, not 10,000.
- Hot keys break sharding. A single Redis shard serves ~100–200K simple ops/sec. A celebrity key at 1M reads/sec needs local in-process caching (L1, 1–5s TTL) or key replication across N shards — sharding alone cannot help a single key.
- Precompute when the read is expensive and the shape is known. Fan-out-on-write feeds, materialized views and denormalized documents turn a 50ms join into a 1ms key lookup — at the cost of write amplification and a backfill problem that someone must own.
The Problem#
Read volume grows faster than write volume and concentrates on a small set of keys. A single Postgres primary on a large instance handles ~20–50K simple indexed reads/sec before CPU or connection limits bite; a feed query that joins follows × posts × likes costs 20–200ms and a few thousand of those per second saturates it. Add a viral post, a marketing push, or a cache flush during deploy and the read path either stays up because it was designed for tails, or it collapses onto the primary and takes writes down with it. The job is to serve the right copy at the right freshness — and to know exactly which reads are not allowed to be stale.
Case Studies That Use This Pattern#
- Distributed Caching — The canonical read-scaling layer: cache-aside, invalidation, hot keys, stampedes
- CDN & Edge Caching — Pushing reads to 200+ PoPs; the edge absorbs 90–99% of static and semi-static traffic
- News Feed — Fan-out-on-write vs fan-out-on-read is a read-scaling decision priced in write amplification
- URL Shortener — 100:1+ read ratio; redirect latency lives or dies on cache hit ratio
- Search Indexing — A secondary read model built from the source of truth, with index lag as the staleness budget
- Leaderboard — Precomputed sorted sets vs on-demand ranking queries
- Replicated Data Store — Replica reads, read-your-writes, and quorum read cost
The Three Intents#
The phrase "scale the reads" hides at least three incompatible goals. Name them and commit before drawing a box.
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Latency for hot, shareable content (feeds, product pages, profiles) | p99 < 50–100ms globally | CDN + cache-aside + precomputation | Stampede on expiry; stale content after edit | Stale ≤ 5–60s acceptable, product-signed |
| Offload / protect the primary (cost and headroom) | Primary CPU < 60%, connection pool < 70% | Read replicas + cache, route by path | Replica lag breaks read-after-write; cache flush sends 20× load to primary | Same as source, lag-bounded |
| Correct reads under concurrency (balances, inventory, ACL checks) | Stale read = money or security incident | Primary reads, lag-aware routing, versioned cache with synchronous invalidation | "Scaled" path silently returns stale value | Read-your-writes or linearizable |
| Expensive query shapes (search, aggregation, analytics) | Query cost 100ms–10s at origin | Secondary read models: search index, OLAP store, materialized views | Index drift from source; backfill takes days | Eventually consistent, drift-monitored |
🎯 Staff Move: "I'll treat this as intent one — latency for shareable content — because that's 95% of the read volume. But I want to carve out the checkout inventory read and the permission check explicitly: those go to the primary or a versioned cache, and I'll pay the cost there on purpose."
The Core Tradeoff#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Cache-aside (Redis/Memcached) | Sub-ms p50, 100K+ ops/sec/node, simple mental model | Stampedes on expiry, stale after write, hot keys on one shard | On-call during flushes; users see stale data for TTL |
| Read replicas | Full SQL on copies, near-zero app change | Replication lag (10ms–minutes), no help for single hot row | Users who just wrote and see old data; support tickets |
| CDN / edge caching | 90–99% offload, 10–30ms global latency | Purge latency (seconds), personalized content can't cache, cache-key explosion | Product (can't personalize), security (cache leaks of private data) |
| Precomputation / materialized views | O(1) read of an expensive shape | Write amplification (1 write → N copies), backfills, schema coupling | Write path latency and storage cost; the team owning backfills |
| Secondary read model (search index, OLAP) | Query shapes the primary can't serve | Index drift, dual-system consistency, rebuild takes hours | Data platform team; users seeing missing/phantom results |
| In-process L1 cache | Zero network hop, absorbs hot keys | N copies × staleness, memory pressure, inconsistent across hosts | Users see different values from different hosts for TTL window |
Staff Default Position#
Budget staleness per read path, then serve each path from the cheapest copy that meets its budget.
The default stack: CDN for anything cacheable without identity, cache-aside with TTL jitter and request coalescing for hot entity reads, read replicas for the long tail of query shapes, and the primary reserved for read-after-write and correctness-critical reads. Precompute only when a read's cost is both high (>20ms) and frequent (>1K/sec) and its shape is stable. Every copy has a named invalidation mechanism and a named owner — "TTL" counts as a mechanism only if product has signed off on the TTL as the staleness bound.
When to Deviate#
- Correctness-critical reads — Balances, inventory at checkout, authorization after revoke. Read from the primary or use versioned keys with synchronous invalidation. The extra 2–5ms is the price of not issuing a refund.
- Write-heavy or uniformly random access — If the working set doesn't fit in cache or access has no skew (hit ratio would be <50%), a cache adds latency and cost without offload. Scale the store horizontally instead (Scaling Writes).
- Extreme hot keys — When one key exceeds a shard's capacity (~100K+ reads/sec), abandon "one copy per key." Replicate the key across N shards or push it into in-process L1 caches with a 1–5s TTL.
- Low scale — Below ~1–2K reads/sec on a well-indexed Postgres, adding a cache adds an invalidation bug surface for no measurable benefit. Add indexes and a connection pooler first.
The L5 → L6 → L7 Contrast#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | "Add Redis in front of the DB" | "What's the staleness budget per read path, and which reads can never be stale?" | "Which read paths across the org share this shape, and should the caching layer be a platform with a staleness contract?" |
| Consistency | TTL of 5 minutes everywhere | Per-path TTL + event-driven invalidation for edited content; read-your-writes via primary routing or session version | Defines an org-wide freshness SLO vocabulary (e.g., F0/F1/F60) that product teams declare per endpoint and platforms enforce |
| Failure | "If Redis goes down we fall through to the DB" | "If Redis goes down the DB sees 20× load; I'll cap fall-through with a concurrency limiter and serve stale" | Designs cache as a load-bearing tier with cell isolation, capacity reserved for cold restart, and a game-day-tested warm-up plan |
| Hot keys | Shard the cache | Detect hot keys (cache.key_qps top-K), L1 cache with 1–5s TTL, request coalescing | Makes hot-key protection a platform default so 40 teams don't each rediscover it during a celebrity event |
| Cost | Not discussed | "Cache hit ratio from 95→99% saves 4 replica nodes" | Prices cache fleet vs replica fleet vs precompute storage in $/month and headcount; retires layers that no longer pay |
| Ownership | Whoever built the service | Service team owns keys and invalidation; SRE owns the cluster | Redraws the contract: platform owns the fleet and client library, product owns key schema, TTL and invalidation correctness |
Why "First move" separates levels
The L5 answer — add Redis — is usually correct and still gets downleveled, because it skips the only question that determines whether the cache is safe: how stale can each read be? The Staff candidate splits the read paths before choosing technology, which lets them say "this path is cacheable for 60s, this one is not cacheable at all." The Principal candidate notices that the same staleness question is being answered ad hoc by every team and turns it into a declared contract.
Why "Failure" separates levels
"Fall through to the DB" assumes the DB can absorb the load the cache was absorbing. At a 95% hit ratio it can't — the origin sees 20× its normal load the instant the cache disappears. The Staff answer bounds fall-through (concurrency limit, serve stale, shed low-priority reads). The Principal answer recognizes the cache is now a load-bearing dependency and funds it like one: cells, reserved warm-up capacity, and rehearsed cold-start.
Why "Ownership" separates levels
Most stale-data incidents are not cache bugs; they are ownership gaps — the team that changed the write path didn't know a cache existed downstream. Staff names the invalidation owner. Principal makes the owner structurally unavoidable: invalidation is emitted from the write path via CDC, not remembered by developers.
The Four Fault Lines#
| # | Fault Line | The Tension |
|---|---|---|
| 1 | Freshness vs Offload | Every second of TTL buys hit ratio and costs correctness |
| 2 | Compute-on-Read vs Compute-on-Write | Pay at read time (latency) or at write time (amplification, backfills) |
| 3 | Shared Cache vs Local Cache | One consistent copy with a network hop, or N fast copies that disagree |
| 4 | Fail-Open vs Fail-Protect on Cache Loss | Fall through to origin (correct, dangerous) or serve stale / shed (safe, wrong) |
Fault Line 1: Freshness vs Offload#
A 5-second TTL on a key read 10K times/sec means ~1 origin query every 5 seconds per key — offload of 99.998%. A 60-second TTL on the same key is barely better for offload but 12× worse for staleness. The curve flattens fast: for hot keys, short TTLs are almost free. For cold keys (read once an hour), no TTL helps — they are always misses. Who pays: users see stale data for the TTL; product must sign off. Staff default: short TTLs (5–60s) on hot data plus event-driven invalidation for anything users edit. Deviate when: data is immutable (content-addressed blobs, versioned assets) — cache forever with versioned keys.
Fault Line 2: Compute-on-Read vs Compute-on-Write#
A feed assembled at read time costs a scatter-gather over followees (50–200ms). A feed precomputed at write time costs one list read (~1ms) — and a celebrity post with 50M followers costs 50M list inserts. Who pays: read-time compute is paid by every reader's latency; write-time compute is paid by the write path, storage, and the team that runs backfills when the schema changes. Staff default: precompute for the median user, compute on read for the heavy tail (the hybrid Twitter-style feed). Deviate when: read shapes change weekly — precomputed views calcify.
Fault Line 3: Shared Cache vs Local Cache#
Redis is one logical copy with a ~0.5–1ms network hop. An in-process cache is ~100ns but there are N copies across N hosts, each expiring independently. Who pays: local caching makes users see values flip between hosts during the TTL window — acceptable for a view count, not for a price. Staff default: two tiers — L1 in-process with 1–5s TTL only for detected hot keys, L2 shared cache for everything else. Deviate when: the fleet is small (<10 hosts) and data is immutable — L1-only is fine.
Fault Line 4: Fail-Open vs Fail-Protect on Cache Loss#
When the cache disappears, origin load multiplies by 1 / (1 − hit_ratio) — 20× at 95%, 100× at 99%. Falling through is "correct" and will take the database down. Who pays: fail-open pays with a total outage; fail-protect pays with stale or missing data for low-priority reads. Staff default: serve stale where possible, cap origin concurrency with a limiter, shed non-critical reads — see Degraded Mode Framework. Deviate when: the origin is genuinely provisioned for full load (rare, and expensive).
The state machine is the talking point: the cache tier has explicit degraded states with entry and exit conditions, and each state has a pre-agreed behavior — nobody decides at 3am whether to serve stale.
Common Interview Mistakes#
| What Candidates Say | What Interviewers Hear | What Staff Engineers Say |
|---|---|---|
| "Add a cache in front of the database" | "One tool for every read path, no staleness analysis" | "Which of these reads can be stale, by how much, and who signed off? Then I'll pick a copy per path." |
| "We'll use read replicas to scale reads" | "Hasn't hit replication lag in production" | "Replicas lag 10–500ms normally and seconds under load. The profile-edit path reads from the primary for 5s after a write." |
| "Set the TTL to 1 hour" | "No idea what the TTL costs the user" | "TTL is a product decision. For prices I'd use event invalidation plus a 60s TTL backstop; for view counts 30s is fine." |
| "If the cache fails we read from the DB" | "Doesn't know the cache is load-bearing" | "At 97% hit ratio, cache loss is 33× origin load. I cap fall-through concurrency and serve stale." |
| "Shard the cache to handle the hot key" | "Doesn't know one key lives on one shard" | "Sharding spreads keys, not load on a key. For a hot key I replicate it or use a 2s L1 cache." |
| "Invalidate the cache after the DB write" | "Race conditions not considered" | "Delete-after-write races with concurrent readers repopulating old values. I'll use versioned keys or leases, and invalidate from CDC." |
Quick Reference#
Staff Sentence Templates#
"Before I pick a cache, I want to split the read paths by staleness budget. [Path A] can be [N] seconds stale and product owns that number; [Path B] cannot be stale at all, so it goes to [the primary / a versioned cache] and I'll pay [X ms] for it."
"Our hit ratio is [H]%, which means losing the cache multiplies origin load by [1/(1−H)]. So the cache is load-bearing, and I'll design the fall-through path as a degraded mode, not an afterthought."
"I'd precompute [the feed / the ranking / the view] because it costs [N ms] and is read [M times/sec]. The price is [K]× write amplification and a backfill problem; [team] owns the backfill runbook."
"Replicas lag [N ms] at p99 under normal load. Any read within [T seconds] of the user's own write routes to the primary, keyed by a session write-timestamp."
Implementation Deep Dive#
1. Cache-Aside with Jitter and Request Coalescing — Redis#
The default. The two production additions most candidates skip are TTL jitter (to desynchronize expiry) and singleflight (so a miss storm becomes one origin query per key per host).
L2_TTL_BASE = 60 # seconds, product-approved staleness budget
inflight = {} # per-process: key -> future
function get(key):
val = redis.GET(key)
if val != null:
metrics.incr("cache.hit", tags=[prefix(key)])
return decode(val)
metrics.incr("cache.miss", tags=[prefix(key)])
# singleflight: only one origin fetch per key per process
if key in inflight:
return inflight[key].await()
fut = new Future()
inflight[key] = fut
try:
row = origin_limiter.run(() => db.query(sqlFor(key))) # bounded concurrency
ttl = L2_TTL_BASE * random.uniform(0.85, 1.15) # ±15% jitter
redis.SET(key, encode(row), EX=ttl)
fut.resolve(row)
return row
except LimiterFull:
stale = redis.GET("stale:" + key) # long-TTL shadow copy
if stale != null:
metrics.incr("cache.served_stale")
return decode(stale)
raise Unavailable
finally:
delete inflight[key]
Why it matters: Without jitter, keys written together (a deploy warm-up, a batch import) expire together. Without coalescing, 5,000 concurrent misses on one key become 5,000 identical DB queries. With both, a hot key costs ~1 origin query per TTL per host.
🎯 Staff Insight: Singleflight is per-process. With 200 app hosts, a stampede still produces up to 200 origin queries per key. That's usually fine; if it isn't, move coalescing into the cache tier with a lease (Memcached-style) or a short Redis
SET NX"refresh lock" so only one host refills.
2. Invalidation Without Races — Versioned Keys and CDC#
"Update DB, then delete cache" has a well-known race: reader A misses, reads old row, writer updates DB and deletes key, reader A writes the old row back into cache. It persists until TTL.
# Write path: bump a version in the same DB transaction
BEGIN;
UPDATE products SET price = $1, version = version + 1 WHERE id = $2 RETURNING version;
COMMIT;
# CDC consumer (Debezium -> Kafka -> invalidator), not application code
on change_event(table="products", id, version):
redis.SET("product:ver:" + id, version) # tiny key, no TTL
redis.DEL("product:" + id + ":v" + (version - 1))
# Read path: key includes the version, so stale fills land on a dead key
function getProduct(id):
v = redis.GET("product:ver:" + id) or db.scalar("SELECT version FROM products WHERE id=$1", id)
return cacheAside("product:" + id + ":v" + v)
Why CDC: invalidation emitted from the database log cannot be forgotten by a developer adding a new write path. The cost is 100ms–2s of invalidation lag (the CDC pipeline) — which becomes your real staleness bound, so measure cdc.invalidation_lag_ms.
3. Read-Your-Writes with Replicas — Session Watermarks#
# On write: remember the primary's WAL position for this session
function onWrite(session):
lsn = db_primary.scalar("SELECT pg_current_wal_lsn()")
session.set("min_lsn", lsn, ttl=10s)
# On read: pick a replica that has replayed at least that far
function routeRead(session):
need = session.get("min_lsn")
if need == null:
return pick(replicas) # no recent write: any replica
for r in replicas_by_lowest_lag():
if r.replay_lsn >= need: # polled every 100ms
return r
metrics.incr("read.routed_to_primary")
return db_primary
Numbers: in steady state >99% of reads go to replicas; only sessions with a write in the last ~10s touch the primary. Alert when read.routed_to_primary exceeds ~5% of reads — it means lag is growing and the primary is about to absorb the read load.
4. Hot Key Defense — Top-K Detection + L1#
# Each app host samples 1% of cache reads into a count-min sketch
every 1s:
hot = sketch.topK(50)
for (key, est_qps) in hot:
if est_qps * 100 > 20_000: # fleet-wide estimate
l1.promote(key, ttl=2s)
function get(key):
if l1.contains(key): return l1.get(key) # ~100ns, no network
return cacheAside(key)
Why 2 seconds: at 1M reads/sec across 500 hosts, a 2s L1 TTL reduces Redis load for that key to ~250 reads/sec (one per host per 2s). The price is up to 2s of cross-host inconsistency — acceptable for a celebrity's follower count, not for a stock price.
Technique Comparison
| Technique | Read Latency | Offload | Staleness | Write Cost | Operational Burden |
|---|---|---|---|---|---|
| Primary read | 1–5ms | 0% | 0 | none | low |
| Read replica | 1–5ms | ~(N−1)/N of reads | 10ms–seconds | replication | medium (lag monitoring) |
| Redis cache-aside | 0.3–1ms | 90–99% | TTL or CDC lag | invalidation | medium |
| L1 in-process | ~0.0001ms | 99.9%+ for hot keys | 1–5s | none | low (but hard to debug) |
| CDN | 10–30ms global | 90–99% of cacheable | TTL / purge 1–10s | purge API | medium |
| Materialized/precomputed | 0.5–2ms | ~100% | pipeline lag | N× write amplification | high (backfills) |
Architecture Diagram#
How to narrate it: each hop down the left-to-right chain is a smaller, fresher, more expensive copy. The CDC stream on the bottom is the single source of invalidation for every copy — the write path never has to know which caches exist. That's the property that keeps stale-data incidents from scaling with the number of teams.
Failure Scenarios#
1. Cache Cluster Restart — 25× Origin Load in 40 Seconds#
A Redis cluster (16 shards, 97% hit ratio, 400K reads/sec) is restarted during a config change. The operator expected a rolling restart; the orchestration restarted all shards within 30 seconds.
t=0 Shards begin restarting. Hit ratio 97% -> 60% -> 5%.
t=+10s Origin read QPS: 12K -> 380K requested. Singleflight collapses to ~140K.
t=+20s Replicas at 100% CPU; p99 read 8ms -> 4s. Connection pools exhausted.
t=+30s App threads blocked on DB; API p99 > 10s; health checks fail; LB ejects hosts.
t=+40s Remaining hosts take ejected hosts' traffic -> cascading ejection.
t=+2min Primary begins serving routed reads (RYW fallback) -> write latency spikes.
t=+6min Manual: enable serve-stale-only mode, limiter to 2K concurrent origin queries.
t=+18min Cache warm; hit ratio back to 95%; limiter relaxed.
Detection: cache.hit_ratio drop >10 points in 1 min; db.replica.cpu >85%; origin_limiter.rejected > 0; read.routed_to_primary > 5%.
Blast radius: every read path and — via RYW fallback — the write path.
Mitigation: origin concurrency limiter (fail fast instead of queueing), serve-stale shadow keys, shed non-critical reads (recommendations, counts).
Prevention: restarts limited to 1 shard at a time with a hit-ratio gate; RYW fallback to primary has its own concurrency cap; cold-start game day each quarter.
Owner: cache platform team owns restart tooling; each service owns its fall-through limiter config.
🎯 Staff Insight: The dangerous line in the timeline is t=+2min — the read incident became a write incident because the read-your-writes fallback had no cap. Every fallback path to the primary needs its own budget.
2. Silent Staleness — Price Change Not Visible for 45 Minutes#
A new bulk-pricing tool writes prices through a stored procedure that the service's "delete cache after update" code never sees. Cached prices stay at the old value until the 1-hour TTL.
t=0 Merchant bulk-updates 30K SKUs via new tool (direct SQL).
t=+1min DB correct. Redis still has old prices (TTL 60min). CDN has old pages.
t=+20min Customers check out at old (lower) prices; orders accepted.
t=+45min Merchant support ticket escalated. Engineer manually flushes product:* keys.
t=+46min Flush causes 15% hit-ratio dip; origin absorbs it (off-peak).
Detection: none fired — the gap. Add cache.staleness_probe (sample 0.1% of reads, compare to primary, emit cache.divergence_rate); alert at >0.5%.
Blast radius: 30K SKUs × 45 min of mispriced orders — a finance issue, not an SRE one.
Mitigation: invalidate from CDC so any write path triggers it; checkout re-reads price from primary.
Prevention: checkout path classified as a zero-staleness read; TTL backstop cut to 5 min for prices.
Owner: catalog team owns invalidation; payments owns the checkout re-read.
3. Replica Lag Spiral — Users "Lose" Their Posts#
A nightly analytics job runs a 20-minute sequential scan on replica 2. Replay on that replica stalls behind conflicting queries; lag climbs to 90 seconds.
t=0 Analytics query starts on replica-2 (should have been on OLAP store).
t=+3min replica-2 lag 5s. Load balancer still sends it 33% of reads.
t=+8min lag 90s. Users who post then refresh see nothing; they post again.
t=+10min Duplicate-post rate 4× baseline; support sees "my post vanished".
t=+14min On-call drains replica-2 from the pool; lag-aware routing added later.
Detection: db.replica.lag_seconds > 5s for 1 min (page); posts.duplicate_rate anomaly.
Blast radius: a third of reads, disproportionately the users who just wrote.
Mitigation: remove lagging replica from rotation automatically when lag > 2s.
Prevention: lag-aware routing (Implementation 3); analytics moved to a dedicated replica excluded from serving; per-query timeouts on serving replicas.
Owner: database platform team owns routing health checks; data team owns the query placement policy.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Cache stampede on hot key expiry | cache.miss spike on one prefix, db.qps step | One key's readers, then origin | Singleflight, jitter, early refresh | Service team |
| Cache cluster loss | cache.hit_ratio < 80% | All read paths, then writes | Origin limiter, serve stale, shed | Cache platform + service |
| Hot key on one shard | redis.shard.cpu skew > 3× median | Every key on that shard | L1 promotion, key replication | Service team (detection is platform) |
| Invalidation miss | cache.divergence_rate > 0.5% | Users of changed entities | CDC-driven invalidation, TTL backstop | Owning domain team |
| Replica lag | db.replica.lag_seconds > 2s | Read-after-write users | Lag-aware routing, auto-drain | DB platform |
| CDN serving private data | Security report; cdn.cache_key audit | Cross-user data leak | Emergency purge, Cache-Control: private | Security + edge team |
| Search index drift | index.doc_count vs source diff > 0.1% | Missing/phantom search results | Reconciliation job, rebuild | Search team |
The Principal Lens#
Why L7 Sees This Problem Differently#
A Staff engineer scales the reads of one system. A Principal engineer notices that the company is running eleven caching clusters, three hand-rolled invalidation schemes, and a CDN configuration nobody fully understands — and that every major stale-data or cache-restart incident in the last year came from the seams between them. At org scale, read scaling stops being a latency technique and becomes a shared dependency with a freshness contract. The questions become: which copies does the org depend on, what freshness does each one promise, who funds the capacity that exists only to survive a cold start, and which layers should be retired because the storage underneath got fast enough.
The Org-Level Fault Line#
One caching platform with a freshness contract vs per-team caches.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Per-team Redis clusters | Autonomy, isolated blast radius, fast to start | 11 clusters × different eviction, TTL, and restart practices; hot-key defense reinvented per team | On-call across 11 teams; security (inconsistent auth/TLS) |
| Central shared cluster | Economies of scale, one expert team | Noisy neighbors; one restart is a company-wide outage | Everyone, simultaneously |
| Platform-owned fleet, cell-per-tenant, shared client library | Consistent defaults (jitter, singleflight, limiter, hot-key L1), isolated blast radius | Platform team becomes a bottleneck for exotic needs | Platform headcount (3–5 engineers) |
The Principal default is the third row: the platform owns the fleet and the client library; product teams own key schemas, TTLs and invalidation correctness; cells are sized so no single restart touches more than one tenant tier.
Cost Model#
Assumptions: 1 KB average value, 95–99% hit ratio, managed Redis at ~$0.20/GB-hour-equivalent node pricing, Postgres replicas on large instances at ~$2.5K/month each, CDN at ~$0.02–0.05/GB egress, fully loaded engineer at ~$25K/month.
| Scale | Read QPS | Cache Fleet | Replicas | CDN | People / On-call | Rough Monthly Total |
|---|---|---|---|---|---|---|
| Startup | 5K | 1 small Redis primary + replica (~$400) | 1 (~$2.5K) | ~$1K | 0.2 FTE, service team on-call | ~$5K + ~$5K people |
| Growth | 150K | 16-shard cluster, 200 GB (~$8K) | 3–4 (~$10K) | ~$15K | 1 FTE cache owner; shared rotation | ~$35K + ~$25K people |
| Large | 3M | 6 cells × 32 shards, 4 TB (~$120K) | 12 across regions (~$30K) | ~$150K | 4-person platform team, dedicated rotation | ~$300K + ~$100K people |
The Principal observation: at the large scale, ~30–40% of cache fleet cost is cold-start headroom — capacity that exists so the origin survives a cache loss. That is an insurance premium, and it should appear on a budget line with an owner, not be discovered during a cost-cutting review and deleted.
The 3-Year Evolution Path#
The Year 3 step is the one most orgs miss: caches are rarely removed. When the underlying store gets fast enough (a DynamoDB/Cassandra read at 2–5ms p99, or a well-indexed Postgres behind a pooler), a cache layer with a 70% hit ratio is pure complexity and invalidation risk.
One-Way Doors vs Two-Way Doors#
| Decision | Door Type | Reversibility Cost |
|---|---|---|
| TTL value for a key prefix | Two-way | Config change, minutes |
| Adding a read replica | Two-way | Hours; remove when idle |
| Exposing stale-tolerant semantics in a public API contract | One-way | Clients build on "eventually consistent"; tightening later is a breaking change |
| Choosing fan-out-on-write for feeds | One-way-ish | Storage layout, backfills and write path all coupled; 2–4 quarters to reverse |
| CDN caching of authenticated responses | One-way (security) | A single leak can't be un-leaked |
| Cache key schema shared across services | One-way-ish | Every consumer must migrate; dual-read period of weeks |
| Choosing CDC as the invalidation source | Two-way (strategic) | Can add other sources later; low regret |
The Standard I'd Write#
RFC: Read-Path Freshness and Caching Standard (v1)
Scope: Every externally reachable read endpoint and every internal read that fans out to >1K QPS.
MUST:
- Each read endpoint declares a freshness tier: F0 (linearizable/primary), F5 (≤5s), F60 (≤60s), F-immutable. Declared in the service's API spec; reviewed by the product owner.
- Services using the shared cache MUST use the platform client (jitter, singleflight, origin concurrency limiter enabled by default).
- Invalidation for mutable entities MUST be driven from CDC or an equivalent log, not from application code alone. TTLs serve as a backstop only.
- Every cache dependency MUST have a documented fall-through budget: max origin QPS the service will send when the cache is unavailable.
- Authenticated responses MUST NOT be cached at the CDN without an approved cache-key review.
SHOULD: Sample 0.1% of F5/F60 reads against the source and publish
cache.divergence_rate; run a cold-start game day twice per year.Exceptions: Filed with the platform team, time-boxed to 2 quarters, owner and expiry recorded.
Success metrics: zero Sev-1s caused by cache restarts; divergence rate <0.1% for F60 tiers; cache fleet cost per 1M reads down 20% YoY; number of distinct cache clusters trending down.
What I'd Tell the VP#
Most of our traffic is people reading things that rarely change, and we serve it from copies. Those copies are why the site is fast and why the database bill is manageable — but twice this year they served wrong prices or took the site down when they restarted. I'm proposing we treat caching as a shared platform with a clear promise about how fresh each page is, rather than eleven teams building their own. It costs roughly two engineers for two quarters, and it pays back by cutting cache-related incidents and letting us retire three clusters. The one thing I need from product leadership is to sign off on how stale each surface is allowed to be.
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Freshness as a contract | "I'd make freshness a declared property of each endpoint, so the cache layer can enforce it and product owns the number." |
| Pricing the insurance | "About a third of this cache fleet is cold-start headroom. I'd put that on the budget explicitly so nobody cuts it in a cost review." |
| Knowing when to retire | "If the store's p99 drops under 5ms, this cache is 70% hit ratio for 100% of the invalidation risk. I'd delete it." |
| Structural invalidation | "Invalidation from CDC means the next team that adds a write path can't forget it. I'm designing out a class of incident, not fixing one." |
| Blast-radius design | "Cells per tenant tier, so a cache restart is a single-cell event, and we game-day it twice a year." |
Staff answers that L7 interviewers find insufficient:
- "I'd add CDC-based invalidation to this service." — Correct for one system; silent on the other ten teams with the same bug.
- "We'll size the cache for 99% hit ratio." — Doesn't ask what the org pays to survive losing it, or who funds that headroom.
- "Redis is the right cache here." — A tool choice without a view on whether the org should run one fleet, eleven, or — in three years — fewer layers.
In the Wild#
Facebook: Scaling Memcache (NSDI 2013)#
Facebook's published memcache architecture serves billions of reads per second from a demand-filled look-aside cache in front of MySQL. Two mechanisms from the paper are directly interview-relevant. Leases let memcache hand a token to the first client that misses on a key; only the lease holder may fill it, which both stops stale sets (the delete/refill race) and collapses thundering herds. Invalidation from the database commit log (a daemon tails MySQL and broadcasts deletes) means application code doesn't have to remember which caches exist. They also used a small "gutter" pool to absorb traffic for failed cache servers instead of falling through to the database.
Staff insight: The two hardest read-scaling problems — stale fills and stampedes — were solved at the cache protocol level, and invalidation was moved to the log. Citing this lets you say "invalidation belongs to the write log, not the developer" with a real reference behind it.
Twitter: Precomputed Home Timelines with a Celebrity Exception#
Twitter has publicly described precomputing home timelines: when a user tweets, the tweet ID is fanned out into the in-memory timeline lists of their followers, so reading a timeline is a single list fetch. For accounts with enormous follower counts, fan-out-on-write is too expensive, so those tweets are merged in at read time.
Staff insight: This is Fault Line 2 resolved as a hybrid, priced explicitly: compute-on-write for the median user, compute-on-read for the heavy tail. It's the cleanest real example of "precompute when the shape is stable and the read is hot."
Netflix: EVCache Replicated Across Zones#
Netflix's EVCache (built on memcached) is publicly documented as storing copies of data in multiple availability zones, with clients reading from the local zone and falling back to another zone on miss or failure. Much of the personalized UI is served from these precomputed, cached results rather than recomputed on request.
Staff insight: Netflix treats the cache as a load-bearing, zone-replicated tier rather than an optimization — the Principal posture. When a zone's cache is lost, reads fail over to another zone's copy instead of stampeding the backing store.
Practice Drill#
Prompt: "Our product catalog page serves 200K reads/sec at peak with a 94% Redis hit ratio. Last Black Friday a Redis failover took the site down for 20 minutes, and merchants complain price edits take up to an hour to show. Fix both."
Staff Answer
Both complaints come from treating the cache as an optimization rather than a load-bearing copy with a freshness contract. First, split the read paths: catalog browse is F60 (product has to sign off on that), the price shown at checkout is F0 and reads from the primary, which is ~3K QPS and easily absorbed. For freshness, move invalidation from application code to CDC: Postgres logical decoding → Kafka → an invalidator that bumps a per-product version key and purges the CDN URL. Staleness drops from "up to TTL (60 min)" to the CDC lag (~1s p99), and I'd cut the TTL backstop to 10 minutes. For resilience: at 94% hit ratio, losing Redis means ~17× origin load (12K → 200K QPS), which the replicas can't take. I'd add (1) a per-host origin concurrency limiter sized so the replica fleet stays under 70% CPU (roughly 15K QPS fleet-wide), (2) a long-TTL shadow copy (stale: keys, 24h) served when the limiter rejects, (3) shedding of recommendations and review counts during cache loss, and (4) a failover procedure that restarts one shard at a time behind a hit-ratio gate. Then game-day it: kill a shard at 50% of peak and verify origin stays under budget. Metrics: cache.hit_ratio per prefix, origin_limiter.rejected, cache.served_stale, cdc.invalidation_lag_ms, cache.divergence_rate.
Why this is L6:
- Classifies reads by staleness budget and carves out the zero-staleness checkout path before touching technology.
- Quantifies fall-through (
1/(1−0.94)≈ 17×) and designs the degraded path with a concrete budget. - Moves invalidation to CDC so the fix outlives the next developer's new write path, and names the metrics that prove it.
What L7 adds:
- Asks whether the other catalog-like surfaces (search, recommendations, seller dashboard) have the same two failure modes and proposes the freshness-tier standard rather than a one-off fix.
- Prices the cold-start headroom (e.g., "keeping replicas at 2× current size costs ~$15K/month — cheaper than one 20-minute Black Friday outage") and puts it on a named budget line.
- Assigns ownership structurally: platform owns the limiter, stale-serving and restart tooling; the catalog team owns key schema and freshness declarations.
Staff Interview Application#
How to Introduce This Pattern#
"Reads here outnumber writes by roughly 100 to 1, so read scaling is where the cost and the risk live. Before I choose caches or replicas, I want to split the read paths by how stale they're allowed to be. Most of this can be 30–60 seconds stale with product sign-off; the checkout and permission reads can't be stale at all, and I'll route those to the primary on purpose."
Lead with the staleness budget, then the copy, then invalidation, then what happens when the copy disappears.
When NOT to Use This Pattern#
- Low read volume: Under ~1–2K QPS on an indexed relational store, add indexes and a connection pooler (PgBouncer) before any cache. The cache would add an invalidation bug surface for single-digit milliseconds.
- Uniform, non-repeating access: If each key is read about once (batch exports, crawls), hit ratio stays near 0%. Scale the store, not a cache.
- Strongly consistent reads dominate: If most reads are F0 (ledgers, inventory at purchase), caching adds risk without offload; scale the primary vertically, partition it, or use a store with strongly consistent reads at scale.
- Write-dominant workloads: Telemetry ingestion, logging — see Scaling Writes.
Follow-Up Questions to Anticipate#
| Interviewer Asks | What They Are Testing | How to Respond |
|---|---|---|
| "How do you invalidate the cache?" | Race awareness | "Not delete-after-write from app code alone — that races with concurrent refills. Versioned keys or leases, with invalidation emitted from CDC; TTL is the backstop." |
| "What if Redis goes down?" | Load-bearing dependency thinking | "Origin load multiplies by 1/(1−hit ratio). I cap origin concurrency, serve stale shadow copies, and shed non-critical reads." |
| "How do you handle a celebrity key?" | Hot-key depth | "Detect via sampled top-K, promote to an in-process L1 with a 2s TTL. At 500 hosts that turns 1M QPS into ~250 QPS on Redis." |
| "User updates their profile and doesn't see it" | Read-your-writes | "Session carries the write's WAL position; reads route to a replica that's caught up or to the primary for ~10s." |
| "Why not just add more replicas?" | Replica limits | "Replicas scale query throughput linearly, but they don't help a single hot row, they add lag, and each one adds replication load on the primary. Past ~5–10 replicas I'd look at caching or partitioning." |
| "When would you precompute?" | Write-amplification pricing | "When the read costs >20ms, runs >1K/sec, and the shape is stable. I'd price the fan-out and name who owns the backfill." |
Evaluation Rubric#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Framing | Picks a cache | Splits reads by staleness budget | Makes freshness a declared, enforced org contract |
| Consistency | TTL only | Versioned keys, CDC invalidation, RYW routing | Invalidation structurally unforgettable across teams |
| Failure | Fall through to DB | Bounded fall-through, serve stale, shed | Cells, funded cold-start headroom, game days |
| Cost | Not discussed | Hit ratio → replica count | $/month per tier, retirement of layers that don't pay |
| Ownership | Implicit | Named service + platform owners | Platform/product contract and an exceptions process |
Strong Hire Signals
| Signal | What It Sounds Like |
|---|---|
| Staleness-first framing | "How stale can this be, and who signed off?" |
| Quantified fall-through | "At 97% hit ratio, cache loss is 33× origin load." |
| Race-aware invalidation | "Delete-then-refill races; I'll version the key." |
| Hot-key specificity | "Sharding doesn't split one key's load." |
Lean No-Hire Signals
| Signal | Why It Misses the Bar |
|---|---|
| "Cache everything, 1-hour TTL" | No staleness analysis; wrong prices are a business incident |
| Replicas with no lag discussion | Will ship read-after-write bugs |
| No answer for cache loss | Doesn't see the cache as load-bearing |
Common False Positives: Deep Redis data-structure knowledge ≠ read-path design. Quoting CDN vendors' PoP counts ≠ knowing what's cacheable. Naming "write-through" ≠ solving invalidation races.
Capacity Planning Quick Reference#
Sizing the Read Path#
origin_qps = total_read_qps × (1 − hit_ratio)
cache_loss_qps = total_read_qps # what origin sees if cache is gone
fallthrough_mult = 1 / (1 − hit_ratio) # 20× at 95%, 100× at 99%
replicas_needed = ceil(origin_qps / per_replica_qps / 0.6) # 60% CPU target
cache_memory = hot_keys × avg_value_bytes × 1.3 # ~30% overhead
l1_backend_qps = hosts / l1_ttl_seconds # per hot key
Key Numbers Worth Memorizing#
| Number | Context |
|---|---|
| 0.3–1 ms | Redis GET p50 in the same AZ |
| ~100–200K ops/sec | Simple ops per Redis shard (single-threaded command execution) |
| ~100 ns | In-process L1 cache lookup |
| 20–50K QPS | Simple indexed reads on one large Postgres node |
| 10–500 ms | Typical async replica lag; seconds under write bursts |
| 90–99% | CDN hit ratio for static/semi-static content |
| 1 s | Typical CDC invalidation lag at p99 |
| ±10–20% | TTL jitter to desynchronize expiry |
| 5× / 10× | Origin reduction going 95→99% / 99→99.9% hit ratio |
| 1–5 s | L1 TTL for hot keys |
Common Pitfalls Checklist#
- Every read path has a declared staleness budget with a product owner
- Zero-staleness reads (balances, checkout, ACL) are explicitly excluded from caches
- TTLs are jittered; misses are coalesced
- Invalidation is driven from the write log, with TTL as backstop
- Fall-through to origin has a concurrency cap and a serve-stale fallback
- Replica routing is lag-aware for read-after-write sessions
- Hot-key detection exists before the celebrity event, not after
- CDN never caches authenticated responses without a cache-key review