skip to content

Consumer Fetch and Parallelism Tuning

Tuning consumer fetch sizes, wait times and poll batch size, and scaling instances within a group. Interviewers ask because most 'the consumers are slow' cases come down to fetch or parallelism configuration.

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

questions

5

What is max.poll.records and why might you lower it when a consumer keeps getting kicked out of its group?

level: juniorimportance: must knowfreq 70%

answer

  1. 500 default records per poll()
  2. max.poll.interval.ms = 5 min deadline between polls
  3. slow batch -> rebalance -> CommitFailedException
  4. lower records OR raise interval
  5. network fetch size is separate

basics

~20 s

max.poll.records caps how many records one poll() call returns. If processing each batch takes too long, the consumer misses its deadline and gets removed from the group. Lowering it means smaller batches that process faster, keeping the consumer alive.

solid answer

~40 s

max.poll.records (default 500) limits the number of records a single poll() call hands back to your application. The Kafka consumer must call poll() again within max.poll.interval.ms (default 5 minutes); poll() also sends heartbeats internally in the background thread, but the actual record-processing happens on your thread between polls. If a batch of 500 records takes longer than max.poll.interval.ms to process, the broker's group coordinator assumes the consumer is dead, triggers a rebalance, and reassigns its partitions. Lowering max.poll.records shrinks each batch so a full process-then-poll cycle finishes inside the interval, preventing these spurious rebalances. The alternative is raising max.poll.interval.ms, but that delays genuine failure detection.

go deeper

for a junior

Know the 500 default and that a too-large/slow batch causes the consumer to be removed from the group; lower it to fix.

for a middle

Distinguish max.poll.interval.ms from session.timeout.ms and explain the rebalance + CommitFailedException chain.

for a senior

Reason about the trade-off vs raising the interval, and when to offload processing to workers with pause/resume.

for a principal

Design consistent poll-loop SLAs across services, set fleet-wide defaults, and teach when batch slicing vs interval tuning is correct.

## The poll loop A Kafka consumer is a single-threaded loop: you call `poll(timeout)`, get a batch of records, process them, then call `poll()` again. The library returns at most `max.poll.records` records per call (default **500**). ## Two different timers (commonly confused) 1. **session.timeout.ms** (default 45s in modern clients): how long the coordinator waits for a **heartbeat** before declaring the member dead. Heartbeats are sent by a **background thread**, not by your processing. 2. **max.poll.interval.ms** (default 300000 = 5 min): the maximum time allowed **between two poll() calls**. This is what catches a consumer that is alive (heartbeating) but stuck processing a batch for too long. ## Why a consumer gets 'kicked out' If processing one batch of records takes longer than `max.poll.interval.ms`, the consumer fails to call poll() in time. The coordinator concludes the member is not making progress, removes it from the group, and triggers a **rebalance** — its partitions are reassigned to other members. When the slow consumer finally calls poll() again, it gets a `CommitFailedException` (its commit is rejected because it no longer owns the partitions), and the same records are reprocessed elsewhere. ## Fixing it with max.poll.records If per-record processing is, say, 200ms, then 500 records = 100s. If a downstream slowdown pushes that past 5 minutes, you rebalance. Setting `max.poll.records=100` caps a batch at ~20s, comfortably inside the interval. Smaller batches = more poll() calls = more frequent deadline resets. ## Trade-offs and alternatives - **Raise max.poll.interval.ms** instead: tolerates slow batches but delays detecting a truly hung consumer. - **Offload processing** to a worker pool and decouple it from the poll thread (advanced; you must pause partitions to avoid unbounded buffering). - max.poll.records does **not** affect how much data is fetched over the network — that's governed by fetch.max.bytes / max.partition.fetch.bytes. It only slices the already-fetched buffer into application batches. ## Edge cases - Setting it too low increases per-call overhead and can hurt throughput. - It interacts with auto-commit: offsets are committed for records actually returned, so smaller batches mean finer-grained commit progress.

  • How is max.poll.interval.ms different from session.timeout.ms?
    session.timeout.ms governs the background heartbeat thread (liveness of the connection). max.poll.interval.ms governs progress — the max gap between poll() calls on your processing thread. A consumer can heartbeat fine yet still be evicted for blowing past max.poll.interval.ms while stuck processing a batch.
  • Does lowering max.poll.records reduce network traffic to the broker?
    No. The fetch from the broker is sized by fetch.max.bytes and max.partition.fetch.bytes. max.poll.records only controls how many already-buffered records each poll() hands to your app, so the same data is fetched regardless.

saying these in an interview costs you the question

  • Saying max.poll.records controls how much data is fetched over the network (it slices the local buffer, not the fetch).
  • Confusing max.poll.interval.ms with session.timeout.ms / heartbeats.
  • Claiming heartbeats are sent during record processing on the main thread (they're on a background thread).

context

open as a page

Explain how fetch.min.bytes and fetch.max.wait.ms work together to trade latency for throughput on the consumer.

level: middleimportance: must knowfreq 65%

basics

~20 s

fetch.min.bytes tells the broker not to answer a fetch until it has at least that many bytes ready. fetch.max.wait.ms caps how long it waits for that. Bigger min.bytes = fewer, larger fetches (more throughput, more latency); the wait prevents waiting forever when traffic is low.

open as a page

Why does adding more consumer instances to a group eventually stop improving throughput, and how do partitions set that ceiling?

level: seniorimportance: must knowfreq 68%

basics

~20 s

Within a consumer group, each partition is consumed by exactly one member. So the number of partitions is the maximum useful parallelism. Once you have as many consumers as partitions, extra consumers sit idle and throughput stops rising.

open as a page

What is max.partition.fetch.bytes versus fetch.max.bytes, and what breaks if max.partition.fetch.bytes is smaller than your largest message?

level: middleimportance: should knowfreq 50%

basics

~20 s

fetch.max.bytes caps the total size of a whole fetch response across all partitions; max.partition.fetch.bytes caps how much comes from each single partition. Historically, if a message was bigger than max.partition.fetch.bytes the consumer could stall — modern clients still return that oversized record so the consumer can progress.

open as a page

A consumer group has growing lag but CPU on the instances is low and the network isn't saturated. Walk through how you'd diagnose and which fetch/parallelism configs you'd suspect.

level: seniorimportance: should knowfreq 45%

basics

~20 s

Low CPU + growing lag usually means the consumers are blocked waiting, not working: a slow downstream call, partition skew (one hot partition), too few partitions to parallelize, or fetches that are too small/infrequent. Check per-partition lag, then look at max.poll.records, fetch sizes, and partition assignment.

open as a page