Revision note (2026-09-15). The earlier version of this article shipped a "robust polling loop with backpressure" whose backpressure branch can never run, a catch (CommitFailedException) around commitAsync() that is unreachable, a recommended fetch.min.bytes of 512 KB with no measurement of what it costs, and verification steps with "expected observations" nobody had observed. It never mentioned max.poll.interval.ms, the setting that actually evicts a slow consumer. This version replaces all of that with numbers from a lab: a single Kafka 4.3.1 broker, kafka-clients 4.3.1, eight tests, and the CLI output pasted as printed.
Three numbers called "lag"
Lag for a consumer group is what kafka-consumer-groups.sh --describe prints: per partition, LOG-END-OFFSET minus CURRENT-OFFSET, where CURRENT-OFFSET is the last committed offset of the group. The Admin API gives the same arithmetic with listConsumerGroupOffsets and listOffsets(OffsetSpec.latest()). In the lab, 900 records go into a 3-partition topic, a consumer with max.poll.records=100 consumes 400 and calls commitSync() once:
[lab] admin lag-ff428c83-0 committed=300 logEnd=300 lag=0
[lab] admin lag-ff428c83-1 committed=100 logEnd=300 lag=200
[lab] admin lag-ff428c83-2 committed=0 logEnd=300 lag=300
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
group-542b887e lag-ff428c83 0 300 300 0 - - -
group-542b887e lag-ff428c83 1 100 300 200 - - -
group-542b887e lag-ff428c83 2 0 300 300 - - -
Both say 500, and the test asserts that they agree. Notice the skew: the consumer drained partition 0 before touching partition 2, so the total hides a partition that has not moved at all. Alert on the maximum per partition, not only on the sum.
The third number is the client-side metric records-lag (and records-lag-max) under kafka.consumer:type=consumer-fetch-manager-metrics. The Kafka monitoring reference says it "is based on current offset and not committed offset", and the difference is not academic. In the lab a consumer reads all 500 records of a partition and commits nothing:
[lab] consumed 500, nothing committed: records-lag=0 records-lag-max(window)=0
The fetcher has reached the end, so the client reports zero. A second test consumes 300 records, then commits offset 0 (which is exactly what a crash before the first commit leaves behind):
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID
group-47b2ad2d uncommitted-a6097f01 0 0 300 300 consumer-group-47b2ad2d-9-...
Same consumer, client metric 0, group lag 300. If your dashboard shows only the client metric, a consumer that processes everything and never commits looks healthy until it restarts and reprocesses the whole topic. Watch the group lag from outside the process (the CLI, the Admin API, or an exporter that reads committed offsets), and use the client metric for what it is good at: telling you whether the fetcher keeps up with the broker.
Tested versions: apache/kafka:4.3.1 broker (KRaft, one node), kafka-clients 4.3.1, Amazon Corretto 17.0.14, Gradle 8.8, Docker 27.4.0. Lab: examples/kafka-consumer-lag.
Partitions set the ceiling on parallelism
One partition goes to at most one consumer of a group. Three consumers on a 2-partition topic, polled until the broker reports a stable group with three members:
[lab] 2 partitions, 3 consumers -> assignment sizes [1, 1, 0]
The third consumer does nothing. Adding instances beyond the partition count changes nothing about lag, and you cannot fix a hot partition by adding consumers either; the records for one key land in one partition and one consumer processes them in order. What helps there is a different key or more partitions, and more partitions do not redistribute existing data.
The fetch settings, measured
The earlier version recommended fetch.min.bytes=524288 (512 KB) with fetch.max.wait.ms of 200 to 500 ms as a baseline. The broker then holds every fetch until 512 KB are available or fetch.max.wait.ms expires. With sparse traffic that is a fixed delay per fetch. The lab produces one 100-byte record at a time and measures how long until poll() returns it, five times per setting:
[lab] fetch.min.bytes=524288 fetch.max.wait.ms=500: [508, 478, 483, 477, 480] median=480
[lab] fetch.min.bytes=1 fetch.max.wait.ms=500: [0, 0, 0, 0, 0] median=0
Half a second of added latency for every record when the topic is quiet (values slightly under 500 ms are fetches that were already waiting when the record arrived). The default of fetch.min.bytes=1 returns immediately. Raising fetch.min.bytes is a throughput optimization for a saturated consumer, and the lab did not measure throughput; do not set it as a "baseline" for a consumer whose problem is latency.
max.poll.records (default 500) does not change fetching at all; it only caps how many of the already-fetched records one poll() hands back, which matters for the next point.
The setting the earlier article left out: max.poll.interval.ms
A consumer that takes longer than max.poll.interval.ms (default 300000 ms, 5 minutes) between two poll() calls is considered failed; the client stops heartbeating, the group rebalances, and its partitions move to someone else. If you process synchronously inside the poll loop, this is the timer that runs while you process. The lab sets it to 5000 ms and processes for 7 seconds:
[lab] processing takes 7 s with max.poll.interval.ms=5000 ...
WARN ... consumer poll timeout has expired. This means the time between subsequent calls to poll() was longer than the configured max.poll.interval.ms ...
[lab] commitAsync returned normally; commitSync -> Offset commit cannot be completed since the consumer is not part of an active group for auto partition assignment; it is likely that the consumer was kicked out of the group.
[lab] commitAsync callback received: CommitFailedException: Offset commit cannot be completed since ...
Two things to take from this. First, this is where CommitFailedException comes from: a commit after the consumer was evicted. The work is done, the offsets are not saved, and the next owner of the partition reprocesses the batch. Lag spikes, and the fix is max.poll.records low enough (or max.poll.interval.ms high enough) that one batch always finishes inside the interval.
Second, commitAsync() does not throw it. The earlier code was:
try {
consumer.commitAsync();
} catch (CommitFailedException e) {
consumer.commitSync(); // "fallback"
}
commitAsync() returns before the request is sent; failures reach only the OffsetCommitCallback. The catch block is dead code, and the fallback never runs. If you want to know that a commit failed, pass a callback, or use commitSync() at the points where it matters (before a rebalance, on shutdown).
The backpressure branch that never runs
The earlier loop kept a ConcurrentLinkedQueue, checked processingQueue.size() > MAX_QUEUE_SIZE before each poll(), and called pause() or resume() accordingly. After the poll it drained the queue completely. The lab runs that loop unchanged in structure, with a 1 ms sleep per record, 2000 records, max.poll.records=500 and a threshold of 100:
[lab] processed=2000 pauses=0 resumes=12 maxQueueSizeAtCheck=0 commitAsyncExceptions=0
The queue is always empty when it is checked, so pause() is never called and resume() is called on every iteration. The loop is really "poll, process everything synchronously, commit", which is a fine loop, but it has no backpressure, and the max.poll.interval.ms clock runs through the whole processQueue() call. Real backpressure needs work to leave the poll thread: hand records to a bounded executor, pause() the partitions whose work is queued, keep calling poll() so the consumer stays alive, and resume() when the queue drains. That design is what pause() and resume() exist for; the lab does not implement it, it only shows that the earlier version did not either.
Session and heartbeat settings are protocol-dependent
The earlier version tuned session.timeout.ms=30000 and heartbeat.interval.ms=10000. With the classic group protocol, which is still the client default in 4.3.1 (group.protocol=classic), those settings are accepted. With the KIP-848 protocol they are rejected outright:
[lab] group.protocol=consumer + session.timeout.ms -> heartbeat.interval.ms, session.timeout.ms cannot be set when group.protocol=CONSUMER
Under group.protocol=consumer the broker owns both values (group.consumer.session.timeout.ms, group.consumer.heartbeat.interval.ms). Say which protocol your tuning is for.
A consumer configuration that reflects the measurements
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9095");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "orders");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "200"); // sized so a batch finishes well inside max.poll.interval.ms
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, "120000"); // 2 min, or as long as your slowest batch really takes
// fetch.min.bytes and fetch.max.wait.ms: leave the defaults unless a saturated consumer measurably benefits
And a loop that commits with a callback and keeps the interval honest:
consumer.subscribe(List.of("orders"));
while (running) {
ConsumerRecords<String, byte[]> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, byte[]> record : records) {
process(record); // must finish, for the whole batch, inside max.poll.interval.ms
}
if (!records.isEmpty()) {
consumer.commitAsync((offsets, error) -> {
if (error != null) {
log.warn("commit failed for {}: {}", offsets, error.toString());
}
});
}
}
consumer.commitSync(); // last one synchronous, on shutdown
What this does not cover
- Throughput. The lab measured only the latency effect of
fetch.min.byteson sparse traffic, not records per second under load for any setting. - Burrow, the JMX exporter, Prometheus and Grafana. The client metric was read in-process through
KafkaConsumer.metrics(); an exporter that reads committed offsets from the broker is the right tool for group lag, and none was set up here. - Rebalance duration, static membership and the cooperative assignor over time.
- Multi-broker clusters and network partitions.
Reproduce it
cd examples/kafka-consumer-lag
./start-broker.sh # apache/kafka:4.3.1 on localhost:9095
gradle test --no-daemon # 8 tests, about 40 s
./stop-broker.sh
Sources
- Kafka 4.3: Consumer configs (fetch.min.bytes, fetch.max.wait.ms, max.poll.interval.ms, max.poll.records, group.protocol, session.timeout.ms)
- Kafka 4.3: Monitoring, consumer fetch metrics (records-lag, records-lag-max)
- Kafka 4.3: Consumer rebalance protocol (KIP-848)
- Kafka 4.3: Basic operations, managing consumer groups
