2.6 Technology deep dives
6.1 The timeline fan-out service (write-path materialization)
Problem it solves. Serve a per-user merged, time-ordered feed at interactive latency for 10M concurrent users, where the query is a 200-way merge that would otherwise run 2M times per second.
Why wasn't the alternative enough?
- Query on read (the SQL JOIN): 400M lookups/s. Also unbounded per-user cost — users following 10,000 accounts can't be served.
- Polling every 5 s: re-asks a question whose answer usually hasn't changed; multiplies the above by online-user count.
- A plain cache of the query result: invalidation is the problem — any of 200 followees posting invalidates it, so the hit rate collapses exactly for active users.
How it works internally.
- On each post: look up followers → append a reference (post ID + timestamp, not the post body) to each follower's timeline structure.
- Timeline structure = a bounded, sorted list per user (Redis sorted set / a capped list), holding maybe the newest ~800–1,000 entries, with older reads falling back to the slow path.
- Fan-out runs asynchronously via a queue, so a spike delays delivery rather than rejecting posts.
- Hybrid: celebrity posts are not fanned out. They're stored once and merged at read time. The read path becomes
merge(materialized_timeline, recent_posts_of_followed_celebrities). - Delivery to online clients uses a subscription/push channel, not polling.
Deployment. A queue (Kafka/BullMQ/SQS) between the write API and a pool of fan-out workers; a keyspace-sharded cache cluster holding timelines; a separate WebSocket/push tier holding client subscriptions.
Monitoring.
- Fan-out lag — time from post accepted → present in the last follower's timeline. This is the SLO ("5 seconds") and must be measured at a high percentile, not the mean.
- Queue depth and consumer lag per partition
- Fan-out factor distribution (p50 / p99 / max) — this is how you detect a new celebrity before they hurt you
- Timeline cache hit rate and eviction rate
- Dropped-write counter (deliberate, for the follow-many case) — you must count what you deliberately drop, or you can't tell it from a bug
Scaling. Shard timelines by user ID. Scale fan-out workers horizontally, partitioning the queue by poster ID. The scaling limit is not throughput but skew: one celebrity post is one queue message that expands to 100M writes. That single message must be split into sub-batches or it becomes a stuck partition.
Backup. Timelines are derived data — the backup is the ability to rebuild them from posts + follows. Test that rebuild, because you will need it after a bad deploy corrupts them. The systems of record (posts, follows, users) need real backups.
What actually breaks in production.
- A user crosses the celebrity threshold and nobody notices. The threshold must be dynamic and monitored, not a constant someone picked in 2019.
- The fan-out queue backs up during a spike and never drains, because delivery cost grew with the backlog — a metastable failure.
- Duplicate or missing posts after a worker crashes mid-fan-out: exactly-once semantics is genuinely hard (Ch 12). The usual practical answer is idempotent writes keyed on
(timeline_owner, post_id). - Unfollow/block/delete races — a post lands in a timeline after the user blocked the author, because the fan-out was already in flight. Requires a read-time filter as a backstop; you cannot fix it purely at write time.
- Cache cluster restart → a stampede of timeline rebuilds hits the database simultaneously.
- Deleting a post requires removing it from N million timelines, or filtering at read time. Almost everyone does read-time filtering and accepts the tombstone cost.
6.2 Percentile-tracking / metrics pipelines (HdrHistogram, t-digest, DDSketch, Prometheus)
Problem it solves. Continuously compute p50/p95/p99/p999 over a rolling window, cheaply, across many machines.
Why wasn't the simple approach enough? Keeping every response time in the window and sorting it every minute is correct but too expensive at high request rates — it's O(n) memory per window and O(n log n) per refresh, per instance.
Why weren't averages enough? The mean tells you nothing about how many users experienced a delay, and it's dominated by neither the typical case nor the tail.
How they work internally.
- HdrHistogram — fixed bucket boundaries with configurable relative precision across a huge dynamic range (µs to hours). Recording is O(1) (an array increment). Memory is fixed and known upfront. Histograms are addable, which is the whole point.
- t-digest — adaptive clustering of the distribution with much finer resolution at the extremes (q→0 and q→1) than in the middle, giving very accurate high percentiles in small space. Mergeable.
- DDSketch — buckets with relative-error guarantees (e.g. every quantile accurate to within 1%), which is the guarantee you actually want for latency. Mergeable.
- Prometheus histograms — fixed
lebuckets;histogram_quantile()interpolates. Accuracy depends entirely on whether you chose bucket boundaries near your SLO. Prometheus summaries compute quantiles client-side and are not aggregatable across instances — the classic trap.
Deployment. Record in-process; export raw bucket counts (never precomputed quantiles) to a collector; aggregate by summing buckets; compute quantiles at query time.
Monitoring the monitoring. Metric cardinality (labels × values — the #1 cause of a monitoring system falling over), scrape duration, and samples dropped.
Scaling. Cardinality control is the entire scaling story: never label a metric with user ID, request ID, or a raw URL path.
What actually breaks.
- Averaged percentiles. Someone downsamples a p99 graph from 15 s to 5 min resolution by averaging, or averages p99 across 50 pods. Both are meaningless, and both silently understate the tail.
- Prometheus summaries aggregated across replicas — same bug, wearing a different hat.
- Bucket boundaries that don't bracket the SLO — your SLO is 200 ms and your buckets jump 100 ms → 250 ms → 500 ms, so
histogram_quantileinterpolates linearly across a bucket and reports fiction. - Server-side-only measurement, missing queueing delay and network latency — the system looks healthy while users suffer (see §2.3).
- Coordinated omission — a load generator that waits for a response before sending the next request systematically fails to record the slow period, because it stopped sending requests during it. HdrHistogram has explicit correction for this; most homegrown benchmarks don't, and they under-report tails by orders of magnitude.
6.3 Overload protection (circuit breakers, load shedding, backpressure, token buckets)
Problem it solves. Prevent a system near capacity from entering the self-sustaining metastable failure state where it can't recover even after load drops.
Why weren't plain retries enough? Retries are the mechanism of the failure. Naive retry converts a transient slowdown into a permanent outage by multiplying load exactly when the system can least afford it.
Why wasn't autoscaling enough? Scaling takes tens of seconds to minutes; a retry storm compounds in seconds. And if the bottleneck is a shared database, adding stateless replicas makes it worse.
How they work internally.
- Exponential backoff with jitter — retry after
base × 2^attempt, multiplied by a random factor. Without jitter, all clients that failed at time T retry together at T+1s, T+3s, T+7s: synchronized thundering herds. Full jitter (random(0, base × 2^attempt)) is the standard recommendation. - Circuit breaker — a state machine per dependency: CLOSED (pass through, count failures) → OPEN (fail fast immediately, don't even try, for a cooldown period) → HALF-OPEN (let a trickle through; success closes it, failure reopens). The value is failing fast: it stops the caller's threads from all blocking on a dead dependency.
- Token bucket — a bucket refilled at rate R with capacity B. Each request consumes a token; empty bucket = reject. Allows bursts up to B while bounding sustained rate at R. Used both for rate limiting and, importantly, for budgeting retries (e.g., retries may consume at most 10% of the token budget — so retries can't dominate under widespread failure).
- Load shedding — the server watches a saturation signal (queue depth, or better, queue wait time) and rejects requests above a threshold with 503/429. Rejecting cheaply is the point: a fast rejection costs almost nothing, a slow timeout costs a thread and a connection.
- Backpressure — propagate "slow down" upstream rather than buffering. In practice: bounded queues (an unbounded queue is the bug), TCP flow control, gRPC/HTTP2 flow control,
Retry-Afterheaders, and reactive-streams style demand signalling. - Adaptive concurrency limits (Netflix's
concurrency-limits, TCP-Vegas-style) — infer the optimal in-flight request limit from observed latency gradient, instead of hardcoding a thread-pool size.
Deployment. Circuit breakers and backoff in the client library (or the service mesh — Envoy/Istio outlier detection). Load shedding at the server edge and also inside the service in front of the expensive resource. Rate limits at the API gateway.
Monitoring. Rejected-request rate by reason; circuit-breaker state transitions per dependency; retry rate as a fraction of total requests (the key early-warning signal); queue wait time (better than queue depth); and, crucially, goodput — successful requests per second — not just throughput.
Scaling. Limits must be per-dependency and adaptive. Static thresholds picked once become wrong the moment hardware or traffic mix changes.
What actually breaks.
- Retries at every layer. Client retries 3×, the SDK retries 3×, the mesh retries 3×, the load balancer retries 3× → one user action becomes 81 backend requests. This is the single most common cause of self-inflicted outages.
- Circuit breaker with too coarse a granularity — one shared breaker for a whole service trips because one endpoint is degraded, taking out healthy functionality.
- Unbounded queues anywhere in the path, which convert "reject fast" into "accept and time out later" — the worst of both worlds, since you did the work and the client gave up.
- Load shedding that sheds the wrong things — dropping health checks or control-plane traffic, so the orchestrator kills instances that were merely busy, and capacity drops under load.
- No jitter, producing synchronized retry waves that look like a clean sawtooth on the graph.
- Timeouts longer than the caller's timeout — the downstream keeps working on a request no one is waiting for, burning capacity on garbage.
6.4 Chaos engineering / fault injection
Problem it solves. Fault-tolerance code is the least-exercised code in the system, and many critical bugs are due to poor error handling. Untested failover is not failover; it's a hypothesis.
Why wasn't testing enough? Unit and integration tests exercise the code you thought of. Production failures come from emergent behavior between systems that doesn't occur when each is tested in isolation, and from environmental assumptions that were true until they weren't.
How it works. Form a hypothesis about steady-state behavior (a business metric, e.g. orders/minute), inject a fault in a bounded blast radius, and check whether steady state holds. Fault types, roughly in order of increasing realism: process kill → CPU/memory pressure → latency injection → packet loss → dependency error injection → network partition → zone loss. Latency injection is consistently the highest-value one, because slow is harder than dead.
Deployment. Start in staging, then production with a small blast radius, during business hours, with an owner watching and a kill switch. Run experiments on a schedule so they keep testing the current system, not the one from six months ago.
Monitoring. The steady-state business metric first; then error rates, saturation, and — the actual point of the exercise — whether the alert fired and whether the runbook worked.
What actually breaks. Running chaos experiments without a kill switch or blast-radius limit; running them at 3 a.m. when nobody can respond; testing only instance termination (the easy case) and never latency or partial failure; and organizations that run chaos experiments but don't fix what they find, which converts the practice into theatre.