03 / 05

Why can long message processing interact badly with consumer-group membership?

Difficulty: 6/10
Distributed-system reasoning, Ordering, Exactly-once

Long Processing vs Consumer-Group Membership: The max.poll.interval Problem

Long message processing interacts badly with consumer-group membership because of the max.poll.interval.ms configuration. A consumer in a group must call poll() at least every max.poll.interval.ms (default 5 minutes). If it does not, the group coordinator considers the consumer dead, removes it from the group, and triggers a rebalance. This is a safety mechanism: it prevents a stuck consumer from holding partitions indefinitely. But if the consumer's processing time for a batch exceeds max.poll.interval.ms, the consumer will be removed from the group even though it is healthy and working. The consequence is a rebalance: the partitions are reassigned to other consumers, the slow consumer's work is wasted, and the records it was processing are re-delivered to another consumer. If this happens repeatedly, the group enters a rebalance loop: consumers keep getting removed and reassigned, processing never completes, and lag grows. This is a common cause of instability in consumer groups that do expensive processing. The trade-off is between the timeout and the processing time: a longer max.poll.interval.ms tolerates longer processing but delays the detection of truly stuck consumers; a shorter one detects stuck consumers faster but causes more false rebalances.

The mechanism that causes the problem is the separation between the poll loop and the heartbeat. Heartbeats are sent by a background thread, so the consumer can send heartbeats even while processing. But max.poll.interval.ms is about the poll loop, not the heartbeat. If the application thread is busy processing and does not return to poll(), the coordinator considers the consumer dead. This is why max.poll.interval.ms is sometimes called the "processing timeout." The fix is to ensure that the time between poll() calls is less than max.poll.interval.ms. There are three approaches: reduce max.poll.records so that each batch is smaller and processing completes faster; move processing to a separate thread pool so that the poll loop can continue calling poll() while the workers process; or increase max.poll.interval.ms if the processing time is legitimately long and bounded. The trade-off is between throughput and rebalance risk. Smaller batches reduce the processing time per poll but increase the number of polls and commits; a separate thread pool increases parallelism but complicates offset management and ordering; a longer timeout tolerates longer processing but delays failure detection. Version note: Kafka 2.4+ introduced incremental cooperative rebalancing, which reduces the impact of rebalances, but it does not eliminate the max.poll.interval issue. Static membership (group.instance.id) allows a consumer to rejoin without triggering a rebalance if it restarts quickly, but it does not help if the consumer is removed due to max.poll.interval. Always tune max.poll.records and max.poll.interval.ms together.

A common mistake is to set max.poll.interval.ms to a very high value to avoid rebalances, which means a truly stuck consumer is not detected and the group may not make progress. Another mistake is to process records asynchronously without managing offsets, which can cause data loss or duplicates. A third mistake is to ignore the relationship between max.poll.records and max.poll.interval.ms; if you increase max.poll.records without increasing the interval, you increase the risk of rebalances. The trade-off is between the cost of rebalances and the cost of delayed failure detection. For a consumer with long processing times, the right approach is to reduce max.poll.records and possibly increase max.poll.interval.ms, and to use static membership to reduce the impact of restarts. For a consumer with short processing times, the defaults are usually fine. Version note: Spring Kafka provides a mechanism to pause the container and continue heartbeating, which can help with long processing. Kafka Streams handles this internally with its own threading model. For a custom consumer, you must manage it yourself.

javascript
  1. 1

    A consumer must call poll() within max.poll.interval.ms or it is removed from the group.

  2. 2

    Long processing can exceed the interval and trigger a rebalance.

  3. 3

    Repeated rebalances cause a rebalance loop and growing lag.

  4. 4

    Reduce max.poll.records, use a thread pool, or increase max.poll.interval.ms.

  5. 5

    Setting the interval too high delays detection of truly stuck consumers.

  6. 6

    Static membership reduces the impact of restarts but not max.poll.interval issues.

  7. 7

    Monitor rebalance rate and consumer logs for the reason.

  8. 8

    Spring Kafka and Kafka Streams handle this internally; custom consumers must manage it.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.