What does KafkaConsumer.poll(Duration) actually do, and why must you call it in a continuous loop?
answer
- poll = the single pump
- returns batch + heartbeats + fetch + commit
- Duration = max block, returns early on data
- stop polling -> kicked from group
- never sleep between polls
basics
~20 spoll(Duration) fetches a batch of records the consumer already has buffered (or waits up to the timeout for some), and also drives heartbeats, group coordination, and offset commits. You loop because all consumer work happens during poll.
solid answer
~40 spoll(Duration) is the single pump of the consumer: each call returns a batch of records (a ConsumerRecords) and, as a side effect, sends heartbeats to the group coordinator, handles rebalances, fetches more data, and (with auto-commit) commits offsets. The Duration is the maximum time to block waiting for records if none are buffered; it returns early once data arrives. Because the consumer is single-threaded and does no background fetching of new work beyond prefetch, you must call poll continuously in a loop. If you stop calling poll, heartbeats stop and the broker eventually evicts you from the group. Typical shape: while running, poll, process the returned batch, repeat. You never sleep between polls or do long work outside the loop.
go deeper
Know poll returns a batch and you must call it in a loop; nothing happens if you don't poll.
Articulate the side effects of poll (heartbeat, fetch, commit) and that Duration is a max block, not fixed.
Explain prefetch buffers, the first-poll setup cost, and how poll cadence interacts with group liveness without conflating it with max.poll.interval.ms.
Teach the single-threaded pump model and design processing patterns (pause/resume, bounded batches) that keep poll cadence healthy at scale.
## The poll loop, from first principles A Kafka **consumer** is the client that reads records (messages) from Kafka topics. Records live on the broker in ordered, numbered logs called **partitions**; the position a consumer has read up to is its **offset**. `KafkaConsumer` is **not** multi-threaded internally for application data — almost everything it does happens *inside* a call to `poll(Duration timeout)`. The canonical usage is: ```java while (running) { ConsumerRecords<K,V> records = consumer.poll(Duration.ofMillis(100)); for (ConsumerRecord<K,V> r : records) { process(r); } } ``` ### What one poll() call does 1. **Returns buffered records.** The consumer keeps a buffer of pre-fetched records per partition. `poll` hands back a batch (`ConsumerRecords`) drawn from that buffer. 2. **Sends fetch requests.** If buffers are low, it issues `Fetch` requests to brokers to refill them (this is the *prefetch pipeline*). 3. **Drives group membership.** It sends **heartbeats** to the **group coordinator** (a broker) so the group knows the consumer is alive, and it processes any **rebalance** in progress. 4. **Commits offsets** if `enable.auto.commit=true` (committed on poll based on `auto.commit.interval.ms`). 5. **Blocks up to `timeout`** only if there is nothing to return; it returns *as soon as* records are available, so the Duration is a ceiling, not a fixed wait. ### Why a continuous loop is mandatory Because heartbeats and fetching are piggybacked on `poll`, if you stop calling it (e.g. you do slow work outside the loop, or `Thread.sleep`), you stop heartbeating and stop fetching. The broker will eventually consider the consumer dead and trigger a rebalance, reassigning its partitions to others. (Note: modern clients run heartbeats on a background thread, but liveness is still gated by how often you call poll via `max.poll.interval.ms` — that specific liveness mechanism is covered by the long-processing leaf.) ### Edge cases - `poll(Duration.ZERO)` returns immediately with whatever is buffered (possibly empty) — useful for non-blocking checks. - The very first `poll` after `subscribe` may also do partition assignment and offset lookup, so it can take longer than the timeout suggests for setup work. - An empty `ConsumerRecords` is normal (no data yet); your loop should handle it gracefully.
- What happens if your processing inside the loop takes too long between poll() calls?The consumer can miss its liveness deadline and be evicted from the group, triggering a rebalance — that liveness window is governed by max.poll.interval.ms (covered in the long-processing leaf). Heartbeats themselves run on a background thread in modern clients, but poll cadence still matters.
- Does poll(Duration.ofMillis(100)) always block for 100ms?No. 100ms is the maximum it will wait when no records are buffered. If records are already available it returns immediately, so under steady load poll returns far faster than the timeout.
saying these in an interview costs you the question
- Saying poll always blocks for the full Duration regardless of data.
- Claiming the consumer fetches data on a fully independent background thread with no relation to poll.
- Thinking you can sleep or do long blocking work outside poll without consequence.
- Believing each poll returns exactly one record (it returns a batch).