Distributed systems
Kafka CommitFailedException — commit cannot be completed since the group has already rebalanced
Written and reviewed by Sahil Srivastav
org.apache.kafka.clients.consumer.CommitFailedException: 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. You can address this either by increasing max.poll.interval.ms or by reducing the maximum size of batches returned in poll() with max.poll.records.
at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.sendOffsetCommitRequest(ConsumerCoordinator.java:1220)What this error actually means
Your consumer was evicted from its group while it was still working, and by the time it tried to commit, the partitions it was committing for belonged to somebody else. The broker rejects the commit because accepting it would let a member that no longer owns a partition move that partition’s offset — which would silently skip records for the member that does own it.
The mechanism is a deadline you probably never set. `poll()` does double duty: it fetches records and it renews the member’s lease. Between two `poll()` calls the group coordinator gives you `max.poll.interval.ms` (default five minutes). If your handler loop spends longer than that processing one batch, the coordinator concludes the member is dead, revokes its partitions, and rebalances. Your thread is alive and healthy — it is simply late.
The important part is what has already happened by the time you see the exception. The records in the uncommitted batch were processed. Another member has now been assigned those partitions from the last committed offset, which is *before* that batch. Those records will be processed again. This exception is therefore not just a commit failure — it is the moment duplicate processing entered your system, and the log line that tells you where to go looking for double-charged, double-shipped, or double-emailed records.
Causes, most common first
- 1One slow record inside an otherwise fast batch. The budget is per `poll()` call, not per record. With `max.poll.records` at its default of 500, a single record that makes a synchronous call to a downstream service with no timeout can burn the whole five minutes on its own and take the other 499 down with it. This is by far the most common shape, and it is why raising `max.poll.interval.ms` only changes the size of the blast radius.
- 2Batch size multiplied by per-record latency exceeds the interval. Arithmetic, not a bug: 500 records at 700 ms each is 350 seconds against a 300-second deadline. It never fires in staging because staging never fills a batch. It fires the first time you replay a backlog, where every poll returns a full batch.
- 3Processing performed on the poll thread, with retries inside it. A handler that retries three times with a one-second backoff turns a 200 ms record into a 3.2-second record on the failure path. When a downstream dependency degrades, every record takes the slow path at once and the interval is breached across the entire group simultaneously.
- 4A stop-the-world pause on the consumer JVM. A long full GC, a swapping host, or a container throttled by its CPU quota can stall the poll loop without any slow code. Check GC logs before you touch Kafka configuration — otherwise you will tune a consumer that was never the problem.
- 5Committing after a rebalance you were told about and ignored. With cooperative rebalancing, `onPartitionsRevoked` is your notice that a partition is leaving. Code that carries on processing an in-flight batch after revocation and then commits will fail here even when it is comfortably inside the interval.
When you see it
- The exception appears in bursts, always on the slowest partitions or the largest batches
- Consumer lag climbs in steps, drops, then climbs again, because work is being redone
- Broker logs show the same group entering `PreparingRebalance` repeatedly with no deploy in progress
- Downstream rows carry duplicates whose timestamps are minutes apart, not milliseconds
- Throughput collapses under load rather than degrading — each rebalance stops the whole group
- The first batch after a restart is the one that fails, because it is the biggest
How to diagnose it
Step 1
Confirm it is the poll interval and not the session timeout
These are two different deadlines with two different symptoms. The heartbeat thread keeps the session alive independently of processing, so a breached `session.timeout.ms` points at a network or GC stall, while a breached `max.poll.interval.ms` points at your handler. The client says which one in the log line immediately before the exception — look for `due to consumer poll timeout has expired`.
grep -E 'poll timeout has expired|Attempt to heartbeat failed|Revoke previously assigned partitions' app.log | tail -40Step 2
Measure the actual time between polls
Do not estimate it. The consumer already reports it: `records-lag-max`, `poll-idle-ratio-avg`, and most decisively `time-between-poll-avg` and `time-between-poll-max`. If max is near your interval and avg is far below it, you have a tail-latency problem in a small number of records, not a sizing problem across all of them.
curl -s localhost:8080/actuator/metrics/kafka.consumer.coordinator.time.between.poll.maxStep 3
Find the group state from the broker side
Run this while the failure is happening. A group that sits in `PreparingRebalance` or `CompletingRebalance` for more than a few seconds is thrashing; a group in `Stable` with high lag is merely slow. The two need completely different work.
kafka-consumer-groups.sh --bootstrap-server broker:9092 --describe --group orders --state
kafka-consumer-groups.sh --bootstrap-server broker:9092 --describe --group orders --members --verboseStep 4
Rule out GC and CPU throttling on the consumer
A single ninety-second pause produces exactly this exception with no slow code anywhere. In a container, also check `nr_throttled` — CPU quota starvation stalls the poll loop just as effectively as a collection does.
jcmd <pid> GC.heap_info
cat /sys/fs/cgroup/cpu.stat | grep -E "nr_throttled|throttled_usec"Step 5
Quantify the duplicates you already have
Every occurrence of this exception reprocessed a batch. Count how many records were replayed by comparing the committed offset before the rebalance with the offset the new owner started from, then check the downstream table for the same business key written twice.
The fix
Bound the work inside one poll, rather than extending the deadline. Set `max.poll.records` to a value where the worst-case per-record latency times the batch size still fits comfortably inside `max.poll.interval.ms`. For a handler whose p99 is 800 ms and a five-minute interval, 100 records leaves a 4x margin; 500 does not.
Put a timeout on every downstream call the handler makes. Without one, a single hung HTTP request makes the poll interval unbounded and no batch size is small enough. This is the fix that actually removes the failure mode; batch sizing only reduces its frequency.
If processing is genuinely long — minutes per record, a job dispatch, a model inference — take it off the poll loop. Hand the record to a bounded executor, and call `pause()` on the assigned partitions while the queue is full and `resume()` when it drains, polling throughout so the lease keeps renewing. This is the supported way to process slowly without lying about liveness.
Make the consumer idempotent, because this exception means you have already replayed records at least once. Key each side effect on the business identifier or on `topic-partition-offset`, and enforce uniqueness in the database rather than checking first and writing after — the check-then-write pattern loses to two consumers running concurrently during a rebalance.
Handle revocation explicitly. In `onPartitionsRevoked`, commit what is finished and abandon what is not; do not attempt to commit a batch for partitions you have been told you no longer own. With `CooperativeStickyAssignor` most partitions survive the rebalance, so the cost of doing this correctly is small.
// Fails whenever one record is slow: the whole batch shares one deadline,
// and the commit at the end is rejected once the group has moved on.
while (true) {
var records = consumer.poll(Duration.ofMillis(500)); // up to 500 records
for (var r : records) {
shipmentClient.create(r.value()); // no timeout
}
consumer.commitSync();
}
// Bounded batch, bounded per-call latency, per-partition commits, and a
// side effect that survives the replay this exception guarantees.
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 100);
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
CooperativeStickyAssignor.class.getName());
while (running) {
var records = consumer.poll(Duration.ofMillis(500));
for (var partition : records.partitions()) {
long lastOffset = -1;
for (var r : records.records(partition)) {
shipmentClient.create(r.value()); // client has connect + read timeouts
lastOffset = r.offset();
}
if (lastOffset >= 0) {
consumer.commitSync(Map.of(partition,
new OffsetAndMetadata(lastOffset + 1)));
}
}
}
-- the other half of the fix, enforced where concurrency cannot argue with it
CREATE UNIQUE INDEX shipments_order_key ON shipments (order_id);How to stop it coming back
- Alert on `time-between-poll-max` crossing half of `max.poll.interval.ms`, not on the exception — by the time the exception fires you already have duplicates
- Treat every downstream client used inside a consumer as requiring an explicit connect and read timeout; make it a review rule, because the default in most clients is infinite
- Load-test with full batches. A backlog replay is the realistic worst case and it is trivial to simulate by pausing the consumer for ten minutes
- Keep the group’s partition count and the consumer’s concurrency in a single place, so scaling out does not quietly change per-instance batch arithmetic
- Assert idempotency in tests by delivering the same batch twice and checking row counts, rather than by reasoning about it
Practise this failure in a real repository
The Gronex challenge ships a consumer whose side effects are applied twice when a batch is redelivered, with a test suite that asserts the end state after a deliberate replay. Passing it requires the same two moves this page describes: bound the work and make the effect idempotent at the storage layer.
FAQ
Should I just increase max.poll.interval.ms?
Only as a stopgap, and only once you know the p99 handler latency you are budgeting for. Raising it to an hour means a genuinely hung consumer holds its partitions for an hour before the group recovers, converting a burst of duplicates into a long outage. Bound the batch and the downstream call instead.
Is this the same as the session timeout being exceeded?
No. `session.timeout.ms` is renewed by a background heartbeat thread and is breached by GC pauses, network partitions, or a dead process. `max.poll.interval.ms` is renewed only by calling `poll()` and is breached by slow processing. The wording in the exception and the preceding `poll timeout has expired` line tell you which.
Does commitAsync avoid the problem?
No — it changes where you find out. The commit is still rejected by the coordinator; you just receive the failure in a callback instead of at the call site, and it is easy to log and ignore. The records are replayed either way.
Will enabling exactly-once semantics remove this?
It removes the duplicate *writes back into Kafka* when the whole pipeline is Kafka-to-Kafka inside a transaction. It does nothing for an external side effect such as an HTTP call or an email, which will still happen twice. See the page on exactly-once versus effectively-once for where the guarantee actually stops.
Related
Other errors engineers hit next to this one
- ERROR: current transaction is aborted, commands ignored until end of transaction block
- ERROR: duplicate key value violates unique constraint
- Deep OFFSET pagination getting slower every page
- ERROR: canceling statement due to lock timeout (ALTER TABLE)
- ERROR: canceling statement due to conflict with recovery
- Threads blocked forever on a ReentrantLock with no deadlock reported
- ThreadLocal value leaking across requests on a pooled thread
- Consumer stuck in Object.wait() with work already in the queue