What we learned building on Kafka at scale
3 June 2025 · James Okafor
Two years into running Kafka as the backbone of a latency-sensitive risk pipeline, we've accumulated opinions. Some of what we learned was in the docs. Most of it wasn't.
Partition count is a one-way door
You can increase partitions. You cannot decrease them without significant pain — rebalancing consumer groups, reprocessing offsets, and accepting a window of ordering violations. Start with more partitions than you think you need. We started with 24 and wished we'd started with 96.
Consumer lag is a lagging indicator
By the time consumer lag is visible in your dashboards, you're already behind. The metric we found more useful was fetch-latency-avg on the broker side — it tends to start climbing earlier, before lag becomes dramatic. Alert on it.
"Consumer lag tells you the house is on fire. Fetch latency tells you the smoke detector went off."
Schema registry is non-negotiable at scale
We shipped without a schema registry for the first eight months. We spent four of those months debugging subtle serialisation mismatches between producer and consumer versions deployed on different schedules. The time we lost to that is embarrassing in retrospect. Schema registry from day one.
What we'd tell ourselves on day one
Run a realistic load test before going live. "Realistic" means 3x your expected peak, sustained for 20 minutes. Kafka handles bursts fine — it's sustained high throughput combined with downstream consumer slowness that causes the problems. Most teams only discover this in production.