advanced 3 min answer

At 14:02 a consumer group of twelve workers is draining a partitioned order-events topic at 9,000 messages a second. By 14:07 throughput is zero while every worker is running and its health check passes, and queue depth is climbing. At 14:12 throughput returns for forty seconds and then stops again. The pattern repeats. What failed, and which design decision allowed it?

kafkaconsumer grouprebalancepoison messagemax.poll.interval
Show the full answer Hide the answer

The trigger

At 14:02:10 one message is a bulk import carrying 40,000 order lines, about 38 MB. Worker 7 receives it in a poll batch and spends more than five minutes inside the handler.

Kafka's max.poll.interval.ms defaults to 300,000 ms, and a member that has not called poll() within that window is removed from the group even though its heartbeat thread is alive, its process is running and its health endpoint returns 200. The heartbeat clock is a different clock: session.timeout.ms defaults to 45,000 ms and worker 7 is passing it the whole time. This is why every per-worker signal is green.

Why it propagated to all twelve

At 14:07:10 the coordinator evicts worker 7 and starts a rebalance. Under the default eager assignors every member revokes every partition and stops consuming until the new assignment is agreed, so one slow member produces a group-wide stall rather than the loss of one twelfth of capacity.

Worker 7 never committed the offset, so the message is redelivered. Worker 3 takes the partition, clears the small messages behind it for forty seconds, reaches the 38 MB message and begins the same five-minute handler. The 14:12 blip is that forty seconds. At 14:17 worker 3 is evicted and the cycle repeats. The poison message tours the group, and each visit costs a full rebalance.

Why detection lagged

Consumer lag cannot tell "too slow" from "stopped", because both look like a rising number. Two signals name the cause and neither is on a default dashboard: the group's rebalance rate, which should be near zero and is now four an hour, and time since last successful poll per member, which pins the stall to one partition and one offset. Add both before you need them; mid-incident nobody thinks to graph a rebalance counter.

Common weak answers

  • "Raise max.poll.interval.ms to an hour." It hides this message and blinds the only mechanism that detects a genuinely hung consumer, so the next deadlocked worker goes unnoticed for an hour instead of five minutes. The eviction clock is a stall detector; you cannot tune away the detector to stop the symptom.
  • "Add consumers." Queue depth with healthy workers reads as a capacity problem, so this is the first instinct. Each new member is another participant in every rebalance, so a group that is spending its time rebalancing gets slower as it grows.
  • "Increase the partition count." It raises the parallelism ceiling and changes nothing here, because the group is not throughput-limited. Applying it also triggers a rebalance, during an incident caused by rebalancing.

The structural fix

Bound the work instead of extending the clock:

  1. Cap work per poll with max.poll.records, so one poll's commitment is bounded whatever arrives.
  2. Take the long handler off the poll thread. Pause the partition, process asynchronously, resume and commit. poll() is then always called on time and the liveness clock keeps its meaning.
  3. Switch to the cooperative sticky assignor (incremental cooperative rebalancing, KIP-429, from Kafka 2.4), so a rebalance revokes only the partitions that move. One evicted member then costs one partition's throughput rather than the group's.
  4. Count delivery attempts and dead-letter after a threshold, so a poison message leaves the group instead of touring it. Without this, the other fixes only make the tour slower.
  5. Limit message size at the producer and use a claim check: write the 38 MB payload to object storage and publish a reference. A queue is a coordination mechanism, not a file transfer.

The general lesson

A timeout that doubles as your liveness detector cannot be relaxed to accommodate slow work; the work has to be made to fit the timeout. And the assignor is a correctness decision rather than a tuning knob: with eager rebalancing, one slow member is a group-wide outage.