skip to content

Consumers and Consumer Groups

The read path: the poll loop, consumer groups and rebalancing, offset commits, and lag. Interviewers spend most of their Kafka time here because nearly every production issue surfaces as a consumer symptom.

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

questions

page 2 of 2

With a consumer using assign() instead of subscribe(), how do you control where consumption starts, since there's no rebalance to trigger position setup?

level: middleimportance: should knowfreq 35%

basics

~10 s

After assign(), the consumer starts from the last committed offset for that group.id if one exists; otherwise auto.offset.reset (earliest/latest) applies. You can override explicitly with seek(), seekToBeginning(), or seekToEnd() before polling.

open as a page

How does pattern subscription (subscribe with a regex) work in Kafka, and what are its operational caveats?

level: middleimportance: should knowfreq 45%

basics

~10 s

subscribe(Pattern) matches topic names against a regex. The consumer auto-discovers new matching topics and rebalances to include them. It still uses group management, just with a dynamic topic set.

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 the difference between position(), committed(), and beginningOffsets()/endOffsets()?

level: middleimportance: should knowfreq 52%

basics

~10 s

position() is where the consumer will read next (in-memory). committed() is the last offset persisted to Kafka for the group. beginningOffsets()/endOffsets() are the oldest and newest (end) offsets currently in each partition.

open as a page

When using the Confluent schema-registry Avro deserializer, what is the difference between SpecificAvro and GenericAvro deserialization, and how does specific.avro.reader control it?

level: seniorimportance: should knowfreq 50%

basics

~10 s

KafkaAvroDeserializer reads a schema-registry ID from each record and fetches the writer schema. With specific.avro.reader=false (default) you get a generic GenericRecord; with true you get generated, typed SpecificRecord classes.

open as a page

What are the generation ID and member.id in the consumer group protocol, and how do they protect against stale/zombie members?

level: seniorimportance: should knowfreq 45%

basics

~20 s

The generation ID is a monotonically increasing counter the coordinator bumps on each successful rebalance; member.id uniquely identifies a member within the group. Requests carrying a stale generation or unknown member.id are rejected, fencing zombies.

open as a page

Why does the replication factor and configuration of the __consumer_offsets topic matter for the Group Coordinator, and what can go wrong if it's misconfigured?

level: seniorimportance: should knowfreq 30%

basics

~10 s

The coordinator's state and committed offsets live in __consumer_offsets. If its replication factor is too low (e.g. 1), losing that broker loses group/offset state and breaks coordinator failover, causing duplicate processing or stuck groups.

open as a page

What are group.min.session.timeout.ms and group.max.session.timeout.ms, and what happens if a consumer requests a session timeout outside them?

level: seniorimportance: should knowfreq 45%

basics

~20 s

They are broker-side limits on the session.timeout.ms a consumer may request: group.min.session.timeout.ms (default 6000 ms) and group.max.session.timeout.ms (default 1800000 ms). A consumer asking for a value outside this range is rejected and cannot join the group.

open as a page

Since KIP-62 split heartbeating from poll(), how do session.timeout.ms and max.poll.interval.ms differ, and why was the split necessary?

level: seniorimportance: should knowfreq 55%

basics

~20 s

session.timeout.ms detects a crashed or network-isolated consumer via missed heartbeats from the background thread. max.poll.interval.ms detects a consumer that is alive but stuck/slow in processing because it hasn't called poll() in time. Before KIP-62 both were conflated into the single session timeout.

open as a page

Why use a dedicated lag monitor like Burrow or Kafka Lag Exporter instead of just the CLI or records-lag-max, and how do their approaches differ?

level: seniorimportance: should knowfreq 45%

basics

~20 s

Dedicated monitors continuously read every group's committed offsets and partition LEOs externally, so they see lag even when consumers are dead, and expose it to Prometheus/alerting. Burrow adds a threshold-free status evaluation; Lag Exporter adds time-based lag estimates.

open as a page

Explain the __consumer_offsets topic: why it's log-compacted, how offsets are keyed, and how the group coordinator uses it.

level: seniorimportance: should knowfreq 45%

basics

~20 s

__consumer_offsets is an internal Kafka topic (50 partitions by default) where consumer commits are stored as messages keyed by (group, topic, partition). It's log-compacted so only the latest committed offset per key is retained. The group coordinator broker reads/writes it to track group progress.

open as a page

How do you commit offsets safely around a consumer group rebalance using ConsumerRebalanceListener?

level: seniorimportance: should knowfreq 60%

basics

~20 s

Register a ConsumerRebalanceListener via subscribe(). In onPartitionsRevoked, commit offsets for partitions you're about to lose before they move to another consumer. onPartitionsAssigned runs when you gain partitions (e.g., to seek to a custom position). This prevents duplicate or lost processing across rebalances.

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

How do you make a consumer start reading from a specific point in time, and what does offsetsForTimes return at the edges?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Use offsetsForTimes(map of partition to timestamp) to look up the first offset whose record timestamp is >= the given time, then seek() to it. If no record has a timestamp at or after the target, it returns null for that partition.

open as a page

How does a ConsumerRebalanceListener's callback behavior differ between eager and cooperative rebalancing, and where should you commit offsets?

level: seniorimportance: should knowfreq 40%

basics

~10 s

ConsumerRebalanceListener has onPartitionsRevoked, onPartitionsAssigned, and onPartitionsLost. Commit offsets in onPartitionsRevoked before giving partitions up. Under eager you're given all partitions; under cooperative you only get the ones actually moving.

open as a page

What is the FencedInstanceIdException, and when does a static consumer get fenced?

level: seniorimportance: should knowfreq 40%

basics

~20 s

If two live consumers join the same group with the same group.instance.id, the broker treats them as duplicates of one static member and fences the older one, throwing FencedInstanceIdException so only one instance keeps consuming each set of partitions.

open as a page

How do you tune static membership for zero-rebalance rolling restarts, and what trade-off does the session timeout introduce?

level: seniorimportance: should knowfreq 45%

basics

~20 s

Give each instance a stable group.instance.id and set session.timeout.ms longer than a single instance's worst-case restart time. Then rolling-restart one instance at a time so each returns inside its window — no rebalances. The trade-off: a longer timeout slows detection of genuinely crashed instances.

open as a page

What is the security risk of deserializing untrusted Kafka payloads, and how do you mitigate it (e.g., with JsonDeserializer trusted packages)?

level: principalimportance: should knowfreq 40%

basics

~20 s

If a deserializer instantiates arbitrary Java classes named in the payload (e.g., via type headers or Java native serialization), an attacker who controls the bytes can trigger remote code execution. Mitigate by restricting allowed types — e.g., Spring's JsonDeserializer trusted.packages allow-list.

open as a page

How do you diagnose partition skew from lag data, and how would you design lag-driven autoscaling for a consumer group?

level: principalimportance: should knowfreq 35%

basics

~20 s

Skew shows as lag concentrated on a few partitions while others sit at zero — usually hot keys or a stuck/slow consumer. For autoscaling, scale consumers on total/time lag, but cap replicas at the partition count, throttle on rebalances, and prefer time-lag with hysteresis to avoid flapping.

open as a page

As a principal engineer, how would you decide between raising max.poll.interval.ms, lowering max.poll.records, and offloading processing for a consumer with slow, variable per-record processing? What are the systemic risks of just cranking the interval up?

level: principalimportance: should knowfreq 40%

basics

~20 s

Match the lever to the cost model: lower max.poll.records if cost scales with count; offload with pause()/resume() if individual records are slow or variable; raise max.poll.interval.ms only for legitimately long, bounded work. Cranking the interval high delays detection of truly stuck consumers, slowing partition reassignment and recovery.

open as a page

You need deterministic, static partition ownership (each service instance always owns the same partitions) without rebalance pauses. How do you design this, and what do you give up?

level: principalimportance: should knowfreq 30%

basics

~20 s

Use assign() to pin a fixed partition set to each instance based on a stable ordinal (e.g. pod index). You avoid rebalances and get deterministic ownership, but you give up automatic load balancing and failover — you must handle scaling and dead-instance takeover yourself.

open as a page

Design how you'd reprocess a topic from a specific point for a running consumer group without losing the current state, and discuss the tradeoffs of the available mechanisms.

level: principalimportance: should knowfreq 34%

basics

~10 s

Use the kafka-consumer-groups CLI --reset-offsets (to earliest, a timestamp, or a specific offset) while the group is stopped, or seek() the partitions in code. auto.offset.reset can't do it because committed offsets already exist.

open as a page

Your consumer group experiences frequent, expensive rebalances during deploys and under load. What assignor and configs would you tune, and why?

level: principalimportance: should knowfreq 38%

basics

~10 s

Switch to CooperativeStickyAssignor so most partitions keep running during rebalances, enable static membership (group.instance.id) so restarts don't trigger rebalances, and tune session.timeout.ms / heartbeat.interval.ms / max.poll.interval.ms so consumers aren't falsely evicted.

open as a page

How does static membership differ from cooperative (incremental) rebalancing, and when would you use each or both?

level: principalimportance: should knowfreq 35%

basics

~20 s

Static membership (group.instance.id) avoids rebalances entirely for transient restarts. Cooperative-sticky rebalancing makes rebalances that DO happen incremental — only affected partitions move instead of stopping the whole group. They solve different problems and are best used together.

open as a page

showing 31–55 of 55