Skip to content
Riadh Mnasri
← Back to blog
7 min read

Kafka: diagnosing consumer lag, methodically

Kafka: diagnosing consumer lag, methodically

An alert fires: "high consumer lag on group billing". Notifications go out twenty minutes late. The most common reflex is to bump the consumer's replica count. It works in one specific case, changes nothing in several others, and makes things worse in at least one. This article offers a way to diagnose lag before touching anything.

What lag measures#

For each partition, Kafka knows two positions: the last offset written by producers (the log end offset) and the last offset committed by the consumer group. Lag is the difference between the two.

partition 0:  [0][1][2][3][4][5][6][7][8][9]
                           ▲              ▲
                  committed offset   log end offset
                           └── lag = 5 ───┘

Two remarks that prevent wrong conclusions.

Lag is read per partition, not as a total. A total lag of 50,000 can be spread evenly over 12 partitions, or concentrated at 49,000 on a single one. Those are two different problems.

Lag in messages says little on its own. 10,000 messages behind on a topic receiving 100,000 per second is a tenth of a second. On a topic receiving 10 per minute, it's two weeks. What matters to users is the delay in time, and above all its trend: stable lag, even high, means the consumer keeps up. Lag that keeps growing means it consumes slower than what comes in, and will never catch up without a change.

Step one: look at the group#

The command shipped with Kafka gives almost everything you need:

bash
kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
  --describe --group billing
TOPIC     PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG     CONSUMER-ID       HOST
invoices  0          184220          184231          11      consumer-1-a3f…   /10.0.1.12
invoices  1          179804          179810          6       consumer-1-a3f…   /10.0.1.12
invoices  2          96115           143902          47787   consumer-2-7c1…   /10.0.1.17
invoices  3          181377          181390          13      consumer-3-e08…   /10.0.1.21

This output already steers the diagnosis. Here, only one partition is behind, and the others keep up. It isn't a global capacity problem.

Three shapes come up almost every time:

What you seeMain lead
All partitions behind, lag growing everywhereProcessing is globally too slow for the throughput
One or two partitions behind, the others up to dateHot key or problematic message on those partitions
Empty CONSUMER-ID column, or one that changes between runsThe group is unstable: rebalances in a loop

Case 1: processing is too slow#

The consumer spends its time processing, and incoming throughput exceeds what it can absorb. Before adding replicas, find out where the time goes. In the vast majority of cases it isn't Kafka, it's what the consumer does with each message: one SQL query per message, a synchronous HTTP call, expensive serialization.

The levers, in the order I try them:

  1. Measure processing time per message. A simple timer around the processing is enough. If a message takes 40 ms and 50 arrive per second per partition, the math is quick: the consumer can't keep up.
  2. Process in batches. max.poll.records (500 by default) sets how many messages each poll returns. One 500-row database insert costs far less than 500 single-row inserts.
  3. Add consumers, but only up to the number of partitions. In a group, a partition is read by a single consumer at a time. With 4 partitions, a fifth consumer sits idle.
  4. Increase the partition count, if you've hit the limit of point 3. It's an operation to plan: for a keyed topic, adding partitions changes each key's target partition, and therefore the processing order during the transition.
Warning

Parallelizing inside a consumer (a thread pool processing the messages of one poll in parallel) raises throughput, but breaks per-partition ordering and complicates offset management: you can only commit an offset once every earlier message is processed. Libraries like Confluent Parallel Consumer do it correctly. Hand-rolling it is a classic source of lost messages.

Case 2: one partition behind#

When a single partition lags, adding consumers does nothing: that partition is still read by only one of them. Two causes dominate.

A hot key. Kafka sends every message with the same key to the same partition, to guarantee their order. If 40% of traffic concerns a single customer (a big account, a default ID like "unknown", or a badly handled null key), its partition gets 40% of the traffic. You check it by counting messages per key on a sample. The fix touches the model: pick a finer key if ordering is only needed at a finer level, or isolate the big customer in a dedicated topic.

A blocking message. A malformed message, one that's too large, or one that triggers an error the consumer retries forever. The partition's lag grows while CURRENT-OFFSET doesn't move at all. That's the clearest sign: a frozen offset. The answer is an explicit error policy: a limited number of attempts, then sending the message to an error topic (dead letter topic) for analysis, and the consumer moves on. Spring Kafka offers this with DefaultErrorHandler and DeadLetterPublishingRecoverer.

Case 3: rebalances in a loop#

This is the case where adding consumers makes things worse. A rebalance redistributes partitions among group members. While it happens, consumption stops (entirely with the historical protocol, partially with the cooperative one). If it happens every few minutes, the group spends more time reorganizing than consuming.

The most common cause is processing that takes longer than max.poll.interval.ms (5 minutes by default). The consumer doesn't call poll in time, the coordinator considers it dead and removes it from the group. Its partitions are reassigned, uncommitted messages are re-read by another consumer, which takes just as long, and the cycle repeats.

poll() → 500 messages × 800 ms = 400 s of processing
       → exceeds max.poll.interval.ms (300 s)
       → the consumer is removed, rebalance
       → the 500 messages are re-read elsewhere → same duration → new rebalance

The signs: logs containing Member ... leaving group or Attempt to heartbeat failed since group is rebalancing, and a CONSUMER-ID column that changes between runs of the command.

The fixes:

  • lower max.poll.records so each batch fits comfortably within the deadline; it's the fastest fix;
  • speed up processing (case 1);
  • raise max.poll.interval.ms, as a last resort, since a genuinely stuck consumer will take longer to detect;
  • switch to cooperative rebalancing (CooperativeStickyAssignor) so only the moved partitions pause, and use static membership (group.instance.id) so a pod restart during a deploy doesn't trigger a rebalance if it comes back in time.

Long GC pauses produce the same symptom through another path: a JVM frozen for several seconds stops sending heartbeats. If rebalances line up with memory spikes, look at the JVM.

Monitor before the alert#

The kafka-consumer-groups.sh command is for diagnosis, not monitoring. For continuous monitoring, two sources:

  • on the consumer side, the records-lag-max JMX metric, exposed automatically by the Java client and picked up by Micrometer in a Spring Boot application;
  • on the cluster side, a Prometheus exporter that computes lag for every group without relying on the consumers, which stays useful when they are all down.

For alerting, I prefer a rule on the trend over a fixed threshold: "this group's lag has been growing for 15 minutes" rather than "lag exceeds 10,000". The fixed threshold wakes you at night for a spike the consumer would have caught up on its own in two minutes, and stays silent on a low-throughput topic where 500 messages mean a day's delay.

The approach in short#

Lag is climbing
 ├─ on every partition?
 │    └─ time per message × throughput > capacity → batches, then replicas (≤ partitions)
 ├─ on one partition?
 │    ├─ frozen offset → blocking message → limited retries + dead letter topic
 │    └─ slowly moving offset → hot key → rethink the key
 └─ changing consumers / rebalance logs?
      └─ processing > max.poll.interval.ms → lower max.poll.records, cooperative, static

Lag isn't a Kafka problem, it's a symptom. Kafka does exactly what it's asked: it keeps messages until someone reads them. The cause is almost always in the consumer or in the choice of keys, and the --describe command, read partition by partition, is enough to know which branch of the tree you're on.