Skip to content
Md. Moudud Hassan
rss

10 min readArchitecture

Adding Consumers Won't Make It Faster

Consumer lag starts climbing. Core banking has slowed from 5,000 messages a second to 2,000, and the gap between what arrives and what you process is widening by the minute.

The instinct is immediate and almost universal: more consumers. Bump concurrency from 3 to 12, scale the deployment from 4 pods to 16, get more hands on the problem.

Lag gets worse. Latency gets worse. Then the consumer group starts rebalancing every few minutes and throughput collapses to well below where it started.

This is one of the few situations in distributed systems where the obvious lever moves the system in the wrong direction, and it is worth understanding exactly why.

Ceiling one: partitions

The cheapest lesson first. A partition is consumed by exactly one member of a consumer group. If a topic has 12 partitions and you start 20 consumers in the same group, 8 of them sit idle — connected, heartbeating, assigned nothing.

So scaling past partition count does nothing at all. Worse than nothing, in fact: every consumer you add is another group member whose arrival and departure triggers a rebalance.

This one is easy to check and easy to fix, though “fix” means repartitioning the topic, which is not free — and if you are keying by account or transaction to preserve ordering, changing partition count changes which key lands where. Plan it; do not do it during an incident.

Ceiling two: you scaled the wrong tier

The deeper problem is that consumers were never the bottleneck.

If core banking sustains 2,000 TPS, then 2,000 TPS is what the system does. Twelve consumers do not make it 6,000. They make twelve consumers compete for a resource that serves 2,000, which changes the failure mode without changing the throughput:

  • Each consumer holds database connections, so the pool drains — the exact connection exhaustion arithmetic, now driven by consumer count.
  • Each concurrent call adds queueing delay downstream, so per-call latency rises for everyone.
  • Rising latency produces timeouts, timeouts produce retries, and retries add load to the thing that was already saturated.

You have converted a clean, visible backlog — messages sitting durably in Kafka, which is exactly where a backlog should sit — into an overloaded downstream system and a retry storm.

Kafka can absorb arrival rate. It cannot raise the throughput of your slowest dependency, and every consumer you add is pressure applied to that dependency rather than capacity added to it.

The spiral: slow processing evicts your consumers

Here is the part that turns a bad afternoon into an outage, and it is a specific, documented Kafka behaviour rather than a vague warning.

A consumer must call poll() at least every max.poll.interval.ms. That is the group coordinator’s liveness check for processing — distinct from the heartbeat thread, which only proves the process is alive. If your handler takes too long on a batch, you miss the deadline, and Kafka concludes the member is stuck.

Kafka’s own CommitFailedException describes it exactly:

Commit cannot be completed since the group has already rebalanced and assigned the partitions to another member. This means that the time between subsequent calls to poll() was longer than the configured max.poll.interval.ms, which typically implies that the poll loop is spending too much time message processing.

Now trace the loop. Downstream slows, so each message takes longer. A batch of 500 that took 20 seconds now takes six minutes and blows past the interval. Kafka evicts the member and rebalances. The partitions move to another consumer, which starts from the last committed offset — so everything processed but not committed is processed again. That reprocessing is additional load on the downstream system that was already too slow, which makes the next batch slower, which triggers the next eviction.

Every turn of that loop adds work. It does not settle on its own, and from the outside it looks like Kafka has broken, when in fact Kafka is doing precisely what it was configured to do.

Tune batch size, not the deadline

Kafka gives you two knobs here, and its documentation is explicit that both exist: raise max.poll.interval.ms, or lower max.poll.records.

Reach for max.poll.records first.

Raising the interval buys time by making Kafka wait longer before concluding a consumer is dead. That delays every rebalance, including the legitimate ones — a genuinely crashed pod now takes minutes to be detected, and its partitions go unconsumed for that whole period. You have traded fast failure detection for slack you will consume again the next time downstream degrades.

Lowering max.poll.records attacks the actual cause. Smaller batches mean each poll cycle finishes in a predictable time regardless of per-message latency, so the deadline stops being at risk:

spring:
  kafka:
    consumer:
      properties:
        # Default is 500. If each message costs ~50ms downstream, a batch of
        # 500 needs 25 seconds — and much longer once downstream degrades.
        max.poll.records: 50
        max.poll.interval.ms: 300000

The arithmetic is worth doing explicitly rather than guessing: worst-case per-message latency × max.poll.records must fit comfortably inside max.poll.interval.ms, with headroom for the degraded case rather than the healthy one. That single calculation prevents most rebalance spirals.

Separate the two knobs

The structural fix is recognising that consumer parallelism and downstream concurrency are different things, and that Kafka only gives you a lever for the first.

Partition count decides how much you can read in parallel. It says nothing about how many concurrent calls your database or core banking system should receive. Coupling them — 12 partitions therefore 12 concurrent downstream calls — means every future repartitioning silently changes the load you apply downstream.

Decouple them explicitly:

@Component
public class PaymentConsumer {

    /**
     * Sized to what core banking can actually absorb from this service, and
     * deliberately NOT derived from partition count. Repartitioning the topic
     * must not silently change the pressure downstream.
     */
    private final Semaphore downstreamPermits = new Semaphore(100);

    @KafkaListener(topics = "payment.events.v1", concurrency = "12")
    public void handle(PaymentEvent event, Acknowledgment ack) throws Exception {
        // Bounded wait: long enough to ride out a blip, short enough that the
        // poll loop still finishes inside max.poll.interval.ms.
        if (!downstreamPermits.tryAcquire(5, TimeUnit.SECONDS)) {
            throw new DownstreamSaturatedException();   // routed to retry topic
        }
        try {
            coreBanking.post(event);
            ack.acknowledge();
        } finally {
            downstreamPermits.release();
        }
    }
}

Twelve consumers now read in parallel while at most 100 calls reach core banking at once. You can repartition the topic without touching downstream pressure, and you can change downstream pressure without repartitioning.

Note the bounded tryAcquire timeout. An unbounded wait would block the poll loop and walk you straight back into the eviction spiral — the protection would have become the cause.

Kafka is already your queue

A common next step is to put a queue between the listener and the downstream call: accept messages quickly, buffer them in memory, drain at a controlled rate. It looks like backpressure.

It is not. It is a second, worse queue.

Kafka is already a durable, replicated, disk-backed buffer designed for exactly this. An in-memory queue in front of a slow dependency moves the backlog from durable storage into volatile heap, and now a pod restart loses it — or the queue grows until the JVM dies. If the queue is bounded, you have merely relocated the blocking; if it is unbounded, you have built an OOM with extra steps.

Let the backlog stay in Kafka. That is what lag is: a number telling you how far behind you are, safely, on disk, replicated three times.

What you should do when the buffer is full is stop reading — which Kafka supports directly. Its own consumer documentation recommends offloading processing to another thread, disabling auto-commit, and pausing the partition so no new records arrive until the work completes. Spring exposes this at the container level:

@Component
@RequiredArgsConstructor
public class ConsumerBackpressure {

    private final KafkaListenerEndpointRegistry registry;
    private final DownstreamHealth health;

    @Scheduled(fixedRate = 5_000)
    public void adjust() {
        MessageListenerContainer container = registry.getListenerContainer("payment-consumer");
        if (container == null) return;

        if (health.isDegraded() && !container.isPauseRequested()) {
            log.warn("Downstream degraded — pausing consumption");
            container.pause();          // finishes in-flight work, stops fetching
        } else if (!health.isDegraded() && container.isPauseRequested()) {
            log.info("Downstream recovered — resuming consumption");
            container.resume();
        }
    }
}

Pausing feels wrong the first time — lag will grow while you are paused. That is the point. Lag is a number; an overwhelmed core banking system is an outage. Trading the first for the second is a good deal, and it is reversible the moment things recover.

Retry off the partition, not on it

One more amplifier. Retrying in place — catch, sleep, try again inside the listener — blocks the partition. Every message behind the failing one waits, ordering stalls, and the sleep counts against your poll deadline.

Route retries to separate topics instead, so the main partition keeps moving:

@RetryableTopic(
    attempts = "4",
    backoff = @Backoff(delay = 1000, multiplier = 3.0, maxDelay = 30000),
    // Saturation is transient; a malformed payload will never succeed.
    exclude = { DeserializationException.class },
    dltStrategy = DltStrategy.FAIL_ON_ERROR)
@KafkaListener(topics = "payment.events.v1")
public void handle(PaymentEvent event) { ... }

Add jitter to the backoff. Without it, everything that failed during a downstream blip retries at the same instant, and you have rebuilt the stampede you were avoiding.

And distinguish failure types deliberately. A timeout deserves a retry; a message that cannot be deserialised will fail identically four more times and belongs in the dead letter topic immediately.

Lag is a symptom, not a diagnosis

Consumer lag is the metric everyone watches and it does not tell you what is wrong. Rising lag has at least three unrelated causes, with opposite remedies:

  • Consumers are genuinely CPU-bound → adding consumers helps, if partitions allow.
  • Downstream is slow → adding consumers hurts, for every reason above.
  • The group is rebalancing → adding consumers makes it much worse.

Distinguishing them takes one extra metric: downstream call latency. If it is flat and lag is climbing, you have a consumer capacity problem — scale out. If it is climbing alongside lag, you have a downstream problem, and scaling out is precisely the wrong response.

Worth alerting on: rebalance frequency (anything routine is a problem), processing time per batch against max.poll.interval.ms, connection pool utilisation, and retry topic depth. Lag alone will lead you to the wrong action about a third of the time.

When adding consumers is right

To be fair to the instinct — sometimes it is correct.

If your consumers are doing real work in-process, deserialising large payloads, running computation, transforming data, and downstream latency is flat while CPU is saturated, then consumers are the bottleneck and adding them is exactly right. Partitions permitting.

The distinction is simply whether the work happens in the consumer or behind it. In-process work scales with consumers. Work behind a shared dependency does not, and no amount of parallelism on your side adds capacity to theirs.

That is also the difference between this and recovering from a backlog after an outage — there the question is how fast to drain, here it is why draining faster is not available to you.

The short version

  1. Consumers beyond partition count do nothing except add rebalances.
  2. If downstream is the constraint, more consumers add pressure, not capacity.
  3. Slow processing triggers rebalances, which cause reprocessing, which adds load — a loop that does not settle.
  4. Lower max.poll.records before raising max.poll.interval.ms. Do the arithmetic: worst-case latency × batch size must fit inside the deadline.
  5. Decouple consumer parallelism from downstream concurrency. They are different numbers.
  6. Do not add an in-memory queue. Kafka is the queue; pause instead.
  7. Retry on separate topics with jittered backoff, never in place.
  8. Watch downstream latency next to lag. Lag alone cannot tell you which problem you have.

The instinct to scale out is correct in most of computing, which is exactly why this case is worth knowing. Here, the honest move is usually to slow down: pause, bound the concurrency, let the backlog sit safely in Kafka, and fix the tier that is actually slow.