skip to content

Describe the consumer poll loop. Why must you call poll() regularly, and what is the role of max.poll.interval.ms?

level: middleimportance: must knowfreq 78%

answer

  1. loop: poll -> process -> commit
  2. poll() = the pulse, drives rebalance
  3. session.timeout = dead consumer (heartbeat)
  4. max.poll.interval = stuck consumer (progress)
  5. exceed it -> rebalance + CommitFailedException

basics

~20 s

After subscribing, you loop calling poll(timeout) to fetch records, process them, then poll again. poll() also drives group membership/heartbeats. If you take too long between polls (longer than max.poll.interval.ms) the broker assumes you're stuck, removes you from the group, and rebalances your partitions to others.

solid answer

~40 s

The standard pattern is: subscribe(topics), then loop { records = poll(Duration); process(records); commit(); }. poll() does more than fetch — it sends and receives the consumer's liveness signals to the group coordinator and runs the rebalance protocol. Two distinct timers govern liveness. session.timeout.ms with heartbeat.interval.ms covers the heartbeat thread: if heartbeats stop (process/JVM dead), the broker evicts you after the session timeout. max.poll.interval.ms covers application progress: it's the max time allowed between successive poll() calls. If your processing of one batch exceeds it, the consumer proactively leaves the group and the coordinator rebalances those partitions, so a subsequent commit fails with CommitFailedException. To avoid this, keep batches small (max.poll.records), make processing fast, or offload work — never block the poll thread for minutes.

code

java · 12 lines
java
consumer.subscribe(List.of("orders"));
try {
    while (running) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
        for (ConsumerRecord<String, String> r : records) {
            handle(r); // keep this fast / bounded
        }
        consumer.commitSync();
    }
} finally {
    consumer.close();
}

go deeper

for a junior

Know the poll/process/commit loop and that you must keep calling poll() or you get kicked out.

for a middle

Distinguish session.timeout.ms (heartbeat/dead) from max.poll.interval.ms (progress/stuck) and explain CommitFailedException.

for a senior

Tune max.poll.records, offload processing, use pause/resume, and reason about at-least-once reprocessing after rebalance.

for a principal

Design processing-time budgets and back-pressure so consumers never trip the interval; standardize on async commit + rebalance listeners across services.

## The loop After you `subscribe(...)` to topics (group-managed), consumption is a loop: ```java consumer.subscribe(List.of("orders")); while (running) { ConsumerRecords<String,String> records = consumer.poll(Duration.ofMillis(100)); for (var record : records) process(record); consumer.commitSync(); } ``` `poll(timeout)` returns the records currently available (up to `max.poll.records`, default 500), waiting at most `timeout` for some to arrive. ## poll() is not just 'fetch' This is the key insight. `poll()` is where the consumer also: (1) joins/rejoins the group and gets partition assignments, (2) runs the rebalance callbacks, (3) updates fetch positions, and (4) in older clients, drove heartbeats. So **you must keep calling poll()** — it's the consumer's pulse, not merely I/O. ## Two liveness timers (don't conflate them) 1. **`session.timeout.ms`** (default 45s) + **`heartbeat.interval.ms`** (default 3s): a background heartbeat thread tells the coordinator 'I'm alive.' If heartbeats stop for `session.timeout.ms` (e.g. the process crashed or paused), the coordinator declares the member dead and rebalances. This detects a *dead consumer*. 2. **`max.poll.interval.ms`** (default 5 minutes): the max wall-clock time allowed *between two poll() calls*. This detects a *live-but-stuck* consumer — heartbeats are still flowing, but the app is spending too long processing one batch and isn't making progress. When exceeded, the consumer **proactively sends a LeaveGroup** and the coordinator rebalances its partitions to other members. Separating these (done in KIP-62) means a slow-processing consumer no longer has to set a huge session timeout; the heartbeat keeps liveness while `max.poll.interval.ms` bounds per-batch processing. ## The classic failure: CommitFailedException If processing a batch exceeds `max.poll.interval.ms`, you've already been kicked out and your partitions reassigned. When you then call `commitSync()`, the broker rejects it: `CommitFailedException` — 'Commit cannot be completed since the group has already rebalanced and assigned the partitions to another member.' The records you processed may be reprocessed by whoever took the partitions (at-least-once). ## How to stay healthy - **Reduce batch size**: lower `max.poll.records` so each poll's batch is processable well within the interval. - **Speed up processing** or move slow work off the poll thread (hand to a worker pool, then commit carefully). - **Raise `max.poll.interval.ms`** only as a last resort if processing is legitimately long; understand it delays detection of truly stuck consumers. - **Call `consumer.pause(partitions)`** if you need more time, keeping the loop spinning (so heartbeats and metadata continue) without fetching more. ## poll timeout vs interval The `Duration` you pass to `poll()` is just how long *that one call* blocks waiting for data — unrelated to `max.poll.interval.ms`. A short poll timeout with fast processing is the healthy shape; the danger is long *processing*, not a long poll timeout.

  • What's the difference between session.timeout.ms and max.poll.interval.ms?
    session.timeout.ms (with the heartbeat thread) detects a crashed/dead consumer when heartbeats stop. max.poll.interval.ms detects a live-but-stuck consumer that is heartbeating but not calling poll() — i.e., processing one batch too long. KIP-62 decoupled them so slow processing doesn't force a huge session timeout.
  • You get CommitFailedException occasionally under load. What's the likely cause and fix?
    Processing a poll batch is exceeding max.poll.interval.ms, so the consumer left the group and its partitions were reassigned before commit. Fix by lowering max.poll.records, speeding up or offloading processing, or (last resort) raising max.poll.interval.ms.
  • If you need to pause processing for a while, how do you keep the consumer in the group?
    Call consumer.pause(partitions) and keep looping poll() — poll() with no fetched records still drives heartbeats and group membership, so you stay alive; resume() when ready. Just don't stop calling poll().

saying these in an interview costs you the question

  • Conflating the poll() timeout argument with max.poll.interval.ms
  • Saying heartbeats are sent by your application code
  • Claiming a slow batch causes session.timeout eviction (it's max.poll.interval)
  • Recommending a giant session.timeout instead of bounding processing

context