1 min lesson
Designing for billions of events per day
Imagine this comes up at work: "Roughly 5 billion events/day arrive. Show how you'd turn that into a sizing target and name the layout decisions that follow." Start with the practical move.
Step 1 of 2
"Billions a day" is a number you should be able to convert to events per second in your head, then turn into partitions, file sizes and consumer counts. Interviewers watch for whether you do the arithmetic or wave at it.
Start with the math, out loud. A day is 86,400 seconds, so one billion events a day averages roughly 11,600 events/sec; five billion is about 58,000/sec. Traffic isn't flat, so apply a peak multiplier - developer usage clusters in working hours, so a 3-5x peak over the daily average is a reasonable planning assumption.
- Daily average
- ~58,000 events/sec (5e9 / 86,400)
- Peak (4x)
- ~230,000 events/sec - size partitions and consumers for this, not the average
- Kafka partitioning
- Enough partitions that per-partition throughput stays under broker/consumer limits, with headroom
- Consumer parallelism
- One consumer per partition ceiling; scale partitions before you can scale consumers
- Target Delta file size
- ~128MB-1GB compacted, not thousands of tiny streaming files
Always size for peak. Sizing for the daily average guarantees you fall over at 9am.