2.10 Worked examples
Query-on-read: 10M online users ÷ 5 s polling = 2M queries/s; × 200 followees = 400M lookups/s.
① The fan-out arithmetic, from scratch. Query-on-read: 10M online users ÷ 5 s polling = 2M queries/s; × 200 followees = 400M lookups/s. Fan-out-on-write: 5,800 posts/s × 200 followers = 1.16M writes/s. Ratio = 400M / 1.16M ≈ 345×. And the write path has a further advantage: it can be queued and delayed during a spike, while the read path cannot — users are waiting.
② Where hybrid fan-out becomes mandatory. A user with 100M followers posts.
Pure write-path: 100M timeline inserts from one event. At 1M writes/s of total system capacity, that single post consumes 100 seconds of the entire cluster's write budget. Read-path merge: 1 stored post + a merge at read time. The threshold isn't a matter of taste — any follower count where followers × post_rate approaches system write capacity forces the read-path treatment.
③ Tail latency amplification, precisely. Each backend call independently exceeds its p99 with probability 0.01.
| Backend calls | P(at least one slow) |
|---|---|
| 1 | 1% |
| 10 | 9.6% |
| 50 | 39.5% |
| 100 | 63.4% |
| 200 | 86.6% |
| At 100 calls, your service's p99 has become your user's median. The only fixes are fewer calls, hedged requests, or a much better per-service p99. |
④ Queueing near capacity. A single-server queue at utilization ρ: average wait ≈ service_time × ρ/(1−ρ).
| ρ | Wait, as a multiple of service time |
|---|---|
| 0.5 | 1× |
| 0.8 | 4× |
| 0.9 | 9× |
| 0.95 | 19× |
| 0.99 | 99× |
| Going from 80% to 95% utilization — a 19% capacity gain — costs you a 5× increase in wait time. This is why "we're only at 90% CPU" is not reassuring. |
⑤ Retry amplification. Client retries 3×, SDK retries 3×, service mesh retries 3×, load balancer retries 3×. Worst case = 3⁴ = 81 backend requests per user action. During a partial outage where most requests fail, your own infrastructure multiplies the load 81-fold at exactly the moment it can least afford it. A retry budget (retries ≤ 10% of traffic) caps this at 1.1×, regardless of how many layers exist.
⑥ Disk failures as normal operation. 10,000 disks at a 3.5% annual failure rate. Failures/year = 350 → ~1 per day. At 100,000 disks: ~10 per day. "In a large-scale system, hardware faults happen often enough that they become part of normal system operation" — this is that sentence as a number.
⑦ Averaging percentiles is meaningless — demonstrated. Two pods, each 1,000 requests. Pod A: 990 requests at 10 ms, 10 at 1,000 ms → p99 = 1,000 ms. Pod B: 1,000 requests at 10 ms → p99 = 10 ms. Averaged p99 = 505 ms. True combined p99 (2,000 requests, the 20th slowest): 10 ms, because only 10 of 2,000 requests exceed it — 0.5%, below the 1% threshold. The average is off by 50×, and in the dangerous direction of neither being right nor conservative.