Everyone had solved this, separately
Kafka was already doing its job. High throughput, fault tolerant, fine at scale — none of that was the problem.
The problem was that every service using it had written its own wrapper. Each team had independently hit the same questions — what do we do when a send fails, how many times do we retry, how does anyone find out — and each had answered them differently, at a different time, under different deadline pressure.
So error handling was inconsistent. Retry logic existed in a dozen near-identical implementations, each with its own quirks. There was no shared way to monitor any of it. And when messages failed, the usual way we found out was a downstream team noticing that something was missing.
During traffic surges this went from untidy to dangerous. Services that were individually well-behaved would collectively overwhelm things downstream, and the first signal was an incident.
What went into the library
A single messaging interface. One API for producing and consuming, with sensible defaults, wrapping the native clients rather than exposing them.
// Simple message production with built-in error handling
kafkaProducer.send(topic, key, value)
.onSuccess(metadata -> log.info("Message sent successfully"))
.onFailure(exception -> alertService.notifyFailure(exception));
Retries that don't make things worse. Naive retry is actively harmful during recovery — every client retries at the same moment, the recovering service gets hit by the full backlog at once, and it falls over again. Exponential backoff with jitter spreads that out.
RetryPolicy policy = RetryPolicy.builder()
.maxAttempts(5)
.exponentialBackoff(Duration.ofSeconds(1), Duration.ofMinutes(5))
.withJitter(0.2)
.build();
Discord webhooks for alerting. This was the least sophisticated piece and the one people reacted to most. When a rate limit trips or a message dies after exhausting its retries, it shows up in the channel the team already has open. The existing alerting was technically better and lived somewhere nobody looked during normal work.
Rate limiting
A token bucket, applied per partition rather than per service. Per-service limiting sounds right and isn't — it lets one hot partition absorb a service's entire budget while the others sit idle. The limit adjusts based on consumer lag, and when the system is under pressure it degrades rather than stopping.
This cut peak-load surges by about 35%.
Observability
The library tracks production latency, consumer lag, retry counts, and rate limit violations, all exposed over JMX so it plugs into the monitoring already in place. Nobody adopts a library that demands a new monitoring stack alongside it.
Where it got to
Surges during peak load down about 35%. P1 incidents from messaging failures down around 60%. Incident response roughly 3x faster, which is almost entirely the Discord alerts — not because the alerts contain better information, but because they arrive somewhere people are already looking. Twelve teams adopted it.
What the adoption curve taught me
Ergonomics drove adoption more than capability. Nobody switched because of the retry algorithm. They switched because the new way was less code than the old way, and the errors told them what to do.
Instrument first. The metrics went in before the first internal release, which meant that when the first real production oddity appeared, the data to diagnose it already existed.
Roll out slowly and boringly. One low-traffic service first, then wider as it held. The edge cases that shook out during that period — mostly around consumer rebalancing — would have been considerably less fun to discover across twelve teams at once.
Still to do
Kafka Streams support, built-in schema validation, and hooks into distributed tracing. All three come from the issue tracker rather than from me, which I take as a reasonable sign the thing is being used.
