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 1 of 2

What are key.deserializer and value.deserializer in a Kafka consumer, and why must they pair with the producer's serializers?

level: juniorimportance: must knowfreq 80%

answer

  1. Brokers store bytes, clients decode
  2. key.deserializer + value.deserializer, no default
  3. Must invert the producer's serializer
  4. Deserializer<T>.deserialize(topic, byte[])
  5. Runs inside poll(); null -> null

basics

~10 s

Kafka stores message keys and values as raw bytes. key.deserializer and value.deserializer are classes that turn those bytes back into objects. They must match the producer's serializers, or the bytes won't decode correctly.

solid answer

~40 s

Kafka brokers store every record's key and value as opaque byte arrays; they never interpret the payload. A consumer must declare key.deserializer and value.deserializer, each implementing org.apache.kafka.common.serialization.Deserializer<T>, to convert bytes back into typed objects. Common built-ins are StringDeserializer, IntegerDeserializer, and ByteArrayDeserializer. The deserializer must be the logical inverse of the producer's serializer: if the producer used StringSerializer (UTF-8), the consumer must use StringDeserializer. A mismatch (e.g., producing Avro bytes but consuming with StringDeserializer) yields garbage or a SerializationException. These are mandatory configs with no default. The deserialize(topic, byte[]) method returns null for null input, which the consumer surfaces as a null key/value.

go deeper

for a junior

Know that messages are bytes and a deserializer turns them into objects; configs are key.deserializer/value.deserializer.

for a middle

Know the built-in deserializers, the pairing-with-serializer rule, and that a mismatch throws SerializationException.

for a senior

Know the Deserializer<T> interface, configure()/isKey, null handling, and that deserialization runs inside poll() and can halt the consumer.

for a principal

Reason about encoding contracts as a cross-team compatibility boundary and where deserialization failures should be handled in the pipeline.

## Why deserializers exist Apache Kafka is **byte-oriented**. A broker receives a record, appends its key bytes and value bytes to a partition log segment, and never inspects them. This decoupling is what lets Kafka carry any payload (JSON, Avro, Protobuf, plain text, images). The cost is that the *client* must agree on encoding. ## Serializer vs deserializer - On the produce side, a `Serializer<T>` converts an object to `byte[]`. - On the consume side, a `Deserializer<T>` converts `byte[]` back to an object. The interface is `org.apache.kafka.common.serialization.Deserializer<T>` with `T deserialize(String topic, byte[] data)` (and an overload taking `Headers`). ## The mandatory configs `key.deserializer` and `value.deserializer` are consumer properties holding the fully-qualified class names. They have **no defaults** — omitting them throws `ConfigException` at consumer construction. Built-ins live in `org.apache.kafka.common.serialization`: - `StringDeserializer` - `IntegerDeserializer` - `LongDeserializer` - `DoubleDeserializer` - `ByteArrayDeserializer` - `ByteBufferDeserializer` - `UUIDDeserializer` ## The pairing rule A deserializer must be the **inverse of the producer's serializer** *for the same field*. Producer `StringSerializer` encodes UTF-8; consumer `StringDeserializer` decodes UTF-8. If a producer wrote Avro-encoded bytes but the consumer uses `StringDeserializer`, you get either mojibake or a `SerializationException` (a subclass of `KafkaException`). Key and value are independent — it's common to use `StringDeserializer` for the key and an Avro deserializer for the value. ## Null handling If `data` is null, the contract is to return null; the `ConsumerRecord` then has a null key or value (used heavily for **tombstones** in compacted topics). ## Configuration injection Deserializers may implement `configure(Map<String,?> configs, boolean isKey)` to read extra properties (e.g., schema-registry URL). The `isKey` flag lets one class behave differently for keys vs values. ## Edge cases Deserialization runs inside `poll()`. An exception thrown there propagates out of `poll()` and, by default, halts the consumer — this is the **'poison pill' problem** addressed by `ErrorHandlingDeserializer`.

  • What happens if you set value.deserializer to a class that can't decode the actual bytes on the topic?
    deserialize() throws a SerializationException inside poll(). By default this propagates and stops the consumer until the offset is skipped or the deserializer is fixed — the classic poison-pill failure.
  • Can the key and value use different deserializers?
    Yes. They are configured independently. A very common pattern is StringDeserializer for the key and an Avro/JSON deserializer for the value.

saying these in an interview costs you the question

  • Claiming the broker validates or understands the payload format (it does not — it only sees bytes)
  • Saying key.deserializer has a default (it has none; omission throws ConfigException)
  • Confusing serializer (produce, object->bytes) with deserializer (consume, bytes->object)

context

open as a page

What is the Group Coordinator in Kafka, and how is the coordinator broker for a particular consumer group chosen?

level: juniorimportance: must knowfreq 70%

basics

~10 s

The Group Coordinator is a broker that manages a consumer group: it handles members joining/leaving, triggers rebalances, and stores offsets. The coordinator is the broker that leads the __consumer_offsets partition the group maps to.

open as a page

What is a consumer heartbeat in Kafka, and what does it accomplish?

level: juniorimportance: must knowfreq 70%

basics

~20 s

A heartbeat is a small periodic signal a consumer sends to the group coordinator broker to prove it is alive and still part of its consumer group. If heartbeats stop, the broker assumes the consumer died and reassigns its partitions.

open as a page

How do you inspect consumer lag from the command line, and what do the columns of kafka-consumer-groups.sh --describe mean?

level: juniorimportance: must knowfreq 70%

basics

~10 s

Run kafka-consumer-groups.sh --bootstrap-server host:port --describe --group <group>. It lists each partition with CURRENT-OFFSET (committed), LOG-END-OFFSET (latest), and LAG = the difference, plus the owning consumer.

open as a page

What is consumer lag in Kafka, and how is it calculated for a single partition?

level: juniorimportance: must knowfreq 80%

basics

~10 s

Consumer lag is how far behind a consumer is on a partition: the latest message offset (log-end-offset) minus the offset the consumer has committed. Lag of 0 means fully caught up.

open as a page

What is max.poll.interval.ms, and what happens if your consumer takes longer than this value to process a batch of records?

level: juniorimportance: must knowfreq 78%

basics

~20 s

max.poll.interval.ms is the maximum time allowed between two poll() calls. If processing a batch takes longer, Kafka assumes the consumer is stuck, kicks it out of the group, and rebalances its partitions to other members.

open as a page

What is the difference between subscribe() and assign() on a Kafka consumer, and when would you choose each?

level: juniorimportance: must knowfreq 75%

basics

~10 s

subscribe() joins a consumer group and lets Kafka assign partitions automatically (dynamic, rebalances on membership change). assign() manually pins specific partitions to the consumer with no group coordination or rebalancing.

open as a page

What is a committed offset in Kafka, and what value does a consumer actually store when it commits?

level: juniorimportance: must knowfreq 80%

basics

~20 s

A committed offset records how far a consumer group has processed in a partition. It stores the offset of the NEXT message to read (last processed + 1), so a restart resumes from there without re-reading.

open as a page

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

level: juniorimportance: must knowfreq 80%

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.

open as a page

What does the consumer config auto.offset.reset control, and what do its values earliest, latest, and none do?

level: juniorimportance: must knowfreq 78%

basics

~10 s

auto.offset.reset decides where a consumer starts reading when there is no valid committed offset for a partition. earliest starts at the oldest message, latest at the newest, and none throws an exception.

open as a page

What is a consumer group rebalance in Kafka, and what triggers one?

level: juniorimportance: must knowfreq 78%

basics

~10 s

A rebalance is when Kafka redistributes a topic's partitions among the consumers in a group. It is triggered when a consumer joins or leaves the group, or when the set of subscribed partitions changes.

open as a page

What is static membership in Kafka consumer groups, and which config enables it?

level: juniorimportance: must knowfreq 60%

basics

~20 s

Static membership gives a consumer a stable identity via the group.instance.id config (KIP-345). With it, a brief consumer restart does not trigger a group rebalance, because the broker recognizes the returning member by its fixed id.

open as a page

Explain client.id, group.id, and isolation.level on a Kafka consumer — what each controls and their valid values.

level: middleimportance: must knowfreq 65%

basics

~10 s

group.id names the consumer group that shares partitions and offsets. client.id is a logical label for the connection used in logs/metrics/quotas. isolation.level (read_uncommitted or read_committed) decides whether the consumer sees records from uncommitted transactions.

open as a page

Walk through the JoinGroup and SyncGroup phases of a rebalance. Who computes the partition assignment, and how does it reach each member?

level: middleimportance: must knowfreq 65%

basics

~20 s

All members send JoinGroup; the coordinator picks one member as group leader and returns the full member list to it. The leader computes assignments and sends them in SyncGroup; the coordinator fans them out to each member in their SyncGroup response.

open as a page

Explain the relationship between heartbeat.interval.ms and session.timeout.ms and how you would tune them.

level: middleimportance: must knowfreq 75%

basics

~20 s

heartbeat.interval.ms is how often the consumer sends heartbeats; session.timeout.ms is how long the coordinator waits without one before declaring the consumer dead. The interval must be well below the timeout — the rule of thumb is heartbeat.interval.ms <= session.timeout.ms / 3.

open as a page

What is the records-lag-max JMX metric, where does it come from, and what are its strengths and blind spots compared to broker-side lag?

level: middleimportance: must knowfreq 60%

basics

~10 s

records-lag-max is a client-side JMX metric exposed by each consumer reporting the maximum lag across the partitions it currently fetches. It's cheap and real-time but only sees assigned partitions of running consumers.

open as a page

Compare automatic offset commits (enable.auto.commit) with manual commits. What are the trade-offs and failure modes?

level: middleimportance: must knowfreq 85%

basics

~20 s

With enable.auto.commit=true, the consumer commits the current position periodically (auto.commit.interval.ms, default 5000ms) during poll(). It's simple but can lose or duplicate records on crash. Manual commits (commitSync/commitAsync after processing) give you control over timing and delivery semantics.

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 seek(), seekToBeginning(), and seekToEnd() — how do they work and when does the new position take effect?

level: middleimportance: must knowfreq 64%

basics

~10 s

These methods manually set where the consumer reads next. seek() jumps to a specific offset; seekToBeginning() to the oldest; seekToEnd() to the newest. The change takes effect on the next poll(), and overrides auto.offset.reset.

open as a page

Compare RangeAssignor, RoundRobinAssignor, and StickyAssignor. When would you pick each?

level: middleimportance: must knowfreq 70%

basics

~10 s

RangeAssignor assigns contiguous partition ranges per topic (can imbalance with multiple topics). RoundRobinAssignor spreads all partitions evenly across consumers. StickyAssignor balances evenly while trying to keep previous assignments stable across rebalances.

open as a page

How does static membership skip a rebalance when a consumer restarts, and what is session.timeout.ms's role?

level: middleimportance: must knowfreq 55%

basics

~20 s

A static member keeps its slot during a brief restart because it doesn't send LeaveGroup and the coordinator hasn't expired its session yet. If it rejoins within session.timeout.ms with the same group.instance.id, it gets the same partitions back — no rebalance.

open as a page

What is a poison-pill record, and how does ErrorHandlingDeserializer solve it?

level: seniorimportance: must knowfreq 70%

basics

~20 s

A poison pill is a record whose bytes can't be deserialized, so it throws on every poll and blocks the consumer forever. ErrorHandlingDeserializer wraps the real deserializer, catches the failure, and hands the error to a recoverer instead of crashing.

open as a page

Walk through how the coordinator detects a dead member via missed heartbeats and what happens next.

level: seniorimportance: must knowfreq 60%

basics

~20 s

The coordinator tracks each member's last heartbeat. If a member sends none for session.timeout.ms, the coordinator declares it dead, removes it from the group, and starts a rebalance that reassigns the dead member's partitions to the remaining consumers.

open as a page

Explain the pause()/resume() backpressure pattern for offloading long-running record processing to worker threads while keeping the consumer in its group.

level: seniorimportance: must knowfreq 56%

basics

~20 s

Hand records to a worker pool for slow processing. To avoid fetching more than you can handle, call consumer.pause() on the assigned partitions and keep calling poll() (which now returns nothing but proves liveness). When workers catch up, call resume(). Commit only completed offsets.

open as a page

Explain the ConsumerRebalanceListener callbacks (onPartitionsRevoked, onPartitionsAssigned, onPartitionsLost) and what you should do in each.

level: seniorimportance: must knowfreq 60%

basics

~20 s

It's a callback you pass to subscribe() so you can react to rebalances: onPartitionsRevoked (commit offsets / flush state before losing partitions), onPartitionsAssigned (init state / seek for new partitions), onPartitionsLost (partitions gone unexpectedly — don't commit).

open as a page

You need at-least-once delivery with no data loss in a consumer. How do you order processing and commits, and how do you bound duplicates?

level: seniorimportance: must knowfreq 70%

basics

~20 s

Disable auto-commit, process the records first, then commit. Never commit before the work is durable. Because a crash between processing and commit replays the batch, make handlers idempotent. For exactly-once, use Kafka transactions or commit offsets atomically with your output to an external store.

open as a page

Explain the difference between eager (stop-the-world) and incremental cooperative rebalancing (KIP-429).

level: seniorimportance: must knowfreq 62%

basics

~20 s

In eager rebalancing every consumer revokes all its partitions and stops processing until a new assignment arrives (stop-the-world). Incremental cooperative rebalancing (KIP-429) only revokes the partitions that actually need to move, so most partitions keep being processed throughout.

open as a page

Describe the consumer group states the coordinator maintains (Empty, PreparingRebalance, CompletingRebalance, Stable, Dead) and the transitions between them.

level: middleimportance: should knowfreq 40%

basics

~10 s

The coordinator tracks a group's lifecycle: Empty (no members), PreparingRebalance (waiting for members to join), CompletingRebalance (awaiting SyncGroup), Stable (assignment done, consuming), and Dead (group removed). Membership changes drive transitions through these states.

open as a page

How does max.poll.interval.ms relate to session.timeout.ms and heartbeat.interval.ms? Why were poll liveness and heartbeat liveness decoupled?

level: middleimportance: should knowfreq 64%

basics

~10 s

session.timeout.ms / heartbeat.interval.ms detect a dead consumer via a background heartbeat thread. max.poll.interval.ms detects a live-but-stuck consumer via the poll() call. They were split (KIP-62) so slow processing doesn't get confused with a crash.

open as a page

A consumer is being evicted because each batch takes too long to process. How does reducing max.poll.records help, and what are the tradeoffs?

level: middleimportance: should knowfreq 58%

basics

~20 s

max.poll.records caps how many records each poll() returns. Lowering it means smaller batches that process faster, so poll() is called again sooner and you stay under max.poll.interval.ms. The tradeoff is more poll round-trips and potentially lower throughput.

open as a page

showing 1–30 of 55