skip to content

Poll Loop and Fetch Mechanics

How poll() actually fetches records: the fetch size and wait knobs, prefetching, and per-partition fetch queues. Interviewers ask because most consumers are tuned, or mis-tuned, through exactly these settings.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

What does KafkaConsumer.poll(Duration) actually do, and why must you call it in a continuous loop?

level: juniorimportance: must knowfreq 80%

answer

  1. poll = the single pump
  2. returns batch + heartbeats + fetch + commit
  3. Duration = max block, returns early on data
  4. stop polling -> kicked from group
  5. never sleep between polls

basics

~20 s

poll(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 s

poll(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

for a junior

Know poll returns a batch and you must call it in a loop; nothing happens if you don't poll.

for a middle

Articulate the side effects of poll (heartbeat, fetch, commit) and that Duration is a max block, not fixed.

for a senior

Explain prefetch buffers, the first-poll setup cost, and how poll cadence interacts with group liveness without conflating it with max.poll.interval.ms.

for a principal

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).

context

open as a page

What does max.poll.records control, and how does it differ from the fetch-size settings?

level: middleimportance: must knowfreq 70%

basics

~20 s

max.poll.records caps how many records a single poll() call returns to your code (default 500). Fetch-size settings (fetch.max.bytes, max.partition.fetch.bytes) control how much data is pulled from brokers over the wire — bytes, not record count.

open as a page

Explain how fetch.min.bytes and fetch.max.wait.ms work together to trade off latency vs throughput in consumer fetches.

level: middleimportance: should knowfreq 55%

basics

~20 s

fetch.min.bytes tells the broker the minimum data to accumulate before answering a fetch; the broker waits up to fetch.max.wait.ms for that much to arrive, then responds anyway. Bigger min.bytes = better throughput but higher latency.

open as a page

What is max.partition.fetch.bytes, how does it relate to fetch.max.bytes, and what failure mode arises if a record exceeds it?

level: seniorimportance: should knowfreq 45%

basics

~20 s

max.partition.fetch.bytes caps bytes returned per partition per fetch (default ~1MB); fetch.max.bytes caps total bytes across all partitions (~50MB). Both are soft limits: to avoid stalling, a broker always returns at least one full record batch even if it exceeds the cap.

open as a page

Describe how the consumer's per-partition fetch buffers and prefetch pipelining work to keep poll() fast and overlap network with processing.

level: seniorimportance: should knowfreq 35%

basics

~20 s

The consumer keeps a buffer of records per partition and issues fetch requests for more data BEFORE the buffer is empty (prefetch). So while you process the current batch, the next fetch is already in flight, and poll() usually returns instantly from the buffer instead of waiting on the network.

open as a page