12.9 Worked examples
Storage = 200e6 × 86,400 × 7 ≈ 121 TB (× replication factor 3 = 363 TB).
① Polling cost vs notification. 1,000 consumers polling a database every 100 ms, of which on average 1% of polls return data. That's 10,000 queries/s, 9,900 of which are wasted. Halve the latency by polling every 50 ms → 20,000 queries/s, 19,800 wasted. The cost is linear in freshness and the yield falls proportionally — §1's exact argument for why push beats poll.
② Log retention as a safety margin. Cluster ingests 200 MB/s; retention 7 days. Storage = 200e6 × 86,400 × 7 ≈ 121 TB (× replication factor 3 = 363 TB). A consumer processing at 180 MB/s that falls behind gains ground at only 20 MB/s. If it was down for 6 hours, it accumulated 200e6 × 21,600 = 4.3 TB of lag, and needs 4.3e12 / 20e6 ≈ 60 hours to catch up. You have 7 days of retention and need 2.5 days to recover — the margin is real but not comfortable. This is why §2.4's "monitor how far behind the head" matters, and why you must know your catch-up rate, not just your lag.
③ Partition count vs parallelism. 12 partitions, 20 consumer instances. 8 consumers sit idle — §2.2's hard ceiling. Now one key carries 40% of traffic: that partition's consumer does 40% of the work while the other 11 share 60%. Adding consumers changes nothing. The only fixes are more partitions and a better key.
④ Window state size. Stream–stream join, 1-hour window, 50,000 events/s, 200 bytes each. State = 50,000 × 3,600 × 200 = 36 GB per side, 72 GB total, held continuously. Extend the window to 24 hours (because some users click days later, §4.3a) → 1.7 TB. This is why §4.2 warns that "large window sizes or high-throughput streams can cause stream processors to keep a lot of temporary state," and why the window length is a capacity decision, not a business decision alone.
⑤ Straggler trade-off. 1-minute tumbling windows; 99.9% of events arrive within 10 s, 0.09% within 5 minutes, 0.01% later.
- Close after 10 s → 0.1% of events dropped; at 50,000 events/s that's 50 events/s discarded, forever, silently unless you count them.
- Close after 5 minutes → 0.01% dropped, but every result is 5 minutes stale.
- Close after 10 s and publish corrections → fresh and eventually accurate, at the cost of downstream consumers having to handle retractions.
⑥ Exactly-once, decomposed. Processing a payment event must: (a) update operator state, (b) write to a database, (c) advance the consumer offset, (d) send a confirmation email.
- (a)+(c) — the framework handles atomically via checkpointing.
- (b) — idempotent if you store the offset alongside the row and check it.
- (d) — nothing the framework can do. The email provider must accept an idempotency key, or you must record "email sent for event X" transactionally with (b). Conclusion: "exactly-once" is always exactly-once within a boundary; every external effect needs its own idempotency story.