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 pageshowhide
explore
- Poll Loop and Fetch Mechanics5 questions
- Heartbeat and Session Liveness5 questions
- Rebalancing and Assignment Strategies5 questions
- Offset Management and Commit Strategies5 questions
- Position, Seek and Offset Reset5 questions
- Consumer Lag and Measurement5 questions
- Static Membership5 questions
- Deserialization and Consumer Configuration5 questions
questions
page 2 of 2With a consumer using assign() instead of subscribe(), how do you control where consumption starts, since there's no rebalance to trigger position setup?
basics
~10 sAfter 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.
How does pattern subscription (subscribe with a regex) work in Kafka, and what are its operational caveats?
basics
~10 ssubscribe(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.
Explain how fetch.min.bytes and fetch.max.wait.ms work together to trade off latency vs throughput in consumer fetches.
basics
~20 sfetch.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.
What is the difference between position(), committed(), and beginningOffsets()/endOffsets()?
basics
~10 sposition() 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.
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?
basics
~10 sKafkaAvroDeserializer 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.
What are the generation ID and member.id in the consumer group protocol, and how do they protect against stale/zombie members?
basics
~20 sThe 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.
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?
basics
~10 sThe 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.
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?
basics
~20 sThey 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.
Since KIP-62 split heartbeating from poll(), how do session.timeout.ms and max.poll.interval.ms differ, and why was the split necessary?
basics
~20 ssession.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.
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?
basics
~20 sDedicated 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.
Explain the __consumer_offsets topic: why it's log-compacted, how offsets are keyed, and how the group coordinator uses it.
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.
How do you commit offsets safely around a consumer group rebalance using ConsumerRebalanceListener?
basics
~20 sRegister 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.
What is max.partition.fetch.bytes, how does it relate to fetch.max.bytes, and what failure mode arises if a record exceeds it?
basics
~20 smax.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.
Describe how the consumer's per-partition fetch buffers and prefetch pipelining work to keep poll() fast and overlap network with processing.
basics
~20 sThe 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.
How do you make a consumer start reading from a specific point in time, and what does offsetsForTimes return at the edges?
basics
~20 sUse 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.
How does a ConsumerRebalanceListener's callback behavior differ between eager and cooperative rebalancing, and where should you commit offsets?
basics
~10 sConsumerRebalanceListener 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.
What is the FencedInstanceIdException, and when does a static consumer get fenced?
basics
~20 sIf 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.
How do you tune static membership for zero-rebalance rolling restarts, and what trade-off does the session timeout introduce?
basics
~20 sGive 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.
What is the security risk of deserializing untrusted Kafka payloads, and how do you mitigate it (e.g., with JsonDeserializer trusted packages)?
basics
~20 sIf 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.
How do you diagnose partition skew from lag data, and how would you design lag-driven autoscaling for a consumer group?
basics
~20 sSkew 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.
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?
basics
~20 sMatch 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.
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?
basics
~20 sUse 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.
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.
basics
~10 sUse 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.
Your consumer group experiences frequent, expensive rebalances during deploys and under load. What assignor and configs would you tune, and why?
basics
~10 sSwitch 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.
How does static membership differ from cooperative (incremental) rebalancing, and when would you use each or both?
basics
~20 sStatic 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.
showing 31–55 of 55