Menu

Kafka course · Lesson 4 of 6

Kafka Consumers: Poll Loop, Offsets, Commits and Lag

How Kafka consumers read data: the poll loop, committed offsets, auto versus manual commits, consumer lag, seeking and replay, deserialisation errors and fetch sizing.

  • Intermediate
  • 23 min read
  • Updated Oct 2026
On this page
  1. Sample data
  2. Consumer API
  3. Worked example: Java consumer with manual commits
  4. Worked example: Python consumer
  5. Pitfalls
  6. In interviews
  7. Poll loop
  8. How it works
  9. Pitfalls
  10. In interviews
  11. Offset management
  12. What happens with no committed offset
  13. Pitfalls
  14. In interviews
  15. Auto vs manual commit
  16. How auto-commit really behaves
  17. Manual commit options
  18. Worked example: commit before versus after processing
  19. Pitfalls
  20. In interviews
  21. Consumer lag
  22. Worked example
  23. Measuring lag in production
  24. Pitfalls
  25. In interviews
  26. Seek and rewind
  27. Worked example: seeking in code
  28. Resetting a group’s offsets from the CLI
  29. Pitfalls
  30. In interviews
  31. Deserialization
  32. Worked example: a poison pill
  33. Patterns for bad records
  34. Pitfalls
  35. In interviews
  36. fetch.min.bytes and fetch.max.bytes
  37. Worked example: catching up on 20,000 records
  38. Tuning guidance
  39. Pitfalls
  40. In interviews
  41. Practice questions
  42. Key takeaways

A producer’s job ends when the broker acknowledges a record. A consumer’s job is harder: it must read, process, write the result somewhere, and record how far it got, in a way that survives crashes and restarts. Getting the order of “process” and “commit” wrong is the most common cause of lost or duplicated data in Kafka pipelines. This lesson covers the consumer from the API to the settings that control how much it fetches.

Sample data

The examples read a three-partition payments topic holding 30 keyed records with amounts 10, 20, … 300 (a sum of 4,650):

bin/kafka-topics.sh --bootstrap-server localhost:19092 --create --topic payments \
  --partitions 3 --replication-factor 3
for i in $(seq 1 30); do echo "acct-$((i%7)):{\"payment_id\":$i,\"amount\":$((i*10))}"; done | \
  bin/kafka-console-producer.sh --bootstrap-server localhost:19092 --topic payments \
  --reader-property parse.key=true --reader-property key.separator=:
bin/kafka-get-offsets.sh --bootstrap-server localhost:19092 --topic payments
payments:0:12
payments:1:4
payments:2:14

The key hash put 12, 4 and 14 records in partitions 0, 1 and 2. Those numbers are the log-end offsets: the offset the next record in each partition will get.

Consumer API

What it is. A consumer is a client that fetches records from partitions. It either subscribes to topics as part of a consumer group (Kafka assigns partitions and tracks progress per group) or assigns itself specific partitions (no group coordination; you manage everything).

Method Use
subscribe(topics) / subscribe(pattern) join the group in group.id; partitions are assigned and rebalanced for you
assign(partitions) read exact partitions without group membership, for tools, replays and tests
poll(Duration) fetch the next records; also drives group membership
commitSync() / commitAsync() record progress for the group
seek, seekToBeginning, seekToEnd, offsetsForTimes move the read position
pause / resume stop fetching from some partitions without leaving the group
wakeup() interrupt a blocking poll() from another thread for shutdown
close() leave the group cleanly (and commit, if auto-commit is on)

The Java KafkaConsumer is not thread-safe: use it from one thread, and use wakeup() (the one thread-safe method) to stop it. Kafka 4.0 removed the old poll(long) method; use poll(Duration), which does not block past its timeout waiting for an assignment.

Worked example: Java consumer with manual commits

import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.*;

public class PaymentConsumer {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:19092");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "payments-loader");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");   // new group: start at the beginning
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");     // commit only after processing
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "10");
        props.put(ConsumerConfig.GROUP_PROTOCOL_CONFIG, "consumer");      // KIP-848 protocol (Kafka 4.0+)

        int processed = 0;
        try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
            consumer.subscribe(List.of("payments"));
            while (processed < 20) {
                ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
                for (ConsumerRecord<String, String> r : records) {
                    // "process" the record: in a real job, write it to the sink here
                    processed++;
                }
                if (!records.isEmpty()) {
                    System.out.printf("polled %2d records from %s%n", records.count(), records.partitions());
                    consumer.commitSync();   // commits the offsets of everything returned by this poll
                }
            }
            Set<TopicPartition> assigned = consumer.assignment();
            Map<TopicPartition, OffsetAndMetadata> committed = consumer.committed(assigned);
            new TreeMap<String, Long>() {{
                committed.forEach((tp, om) -> put(tp.toString(), om == null ? -1 : om.offset()));
            }}.forEach((tp, off) -> System.out.println("committed " + tp + " = " + off));
        }
    }
}
polled 10 records from [payments-2, payments-1]
polled 10 records from [payments-0, payments-2]
committed payments-0 = 2
committed payments-1 = 4
committed payments-2 = 14

The consumer stopped after 20 records. Partitions 1 and 2 are fully read; partition 0 is at offset 2 of 12. A committed offset is the next offset to read, not the last one processed.

Worked example: Python consumer

import json
from confluent_kafka import Consumer, KafkaError, KafkaException

consumer = Consumer({
    "bootstrap.servers": "localhost:19092",
    "group.id": "payments-audit",
    "auto.offset.reset": "earliest",
    "enable.auto.commit": False,
})
consumer.subscribe(["payments"])

total, empty_polls = 0, 0
try:
    while empty_polls < 5:
        msgs = consumer.consume(num_messages=50, timeout=1.0)
        if not msgs:
            empty_polls += 1
            continue
        empty_polls = 0
        for msg in msgs:
            if msg.error():
                if msg.error().code() == KafkaError._PARTITION_EOF:
                    continue
                raise KafkaException(msg.error())
            payment = json.loads(msg.value())
            total += payment["amount"]
        consumer.commit(asynchronous=False)    # commit after the batch is processed
    print("sum of amounts:", total)
    for tp in consumer.committed(consumer.assignment(), timeout=10):
        print(f"committed {tp.topic}[{tp.partition}] = {tp.offset}")
finally:
    consumer.close()
sum of amounts: 4650
committed payments[0] = 12
committed payments[1] = 4
committed payments[2] = 14

A separate group (payments-audit) read all 30 records independently of payments-loader.

Pitfalls

  • Sharing one consumer between threads causes ConcurrentModificationException in Java. Use one consumer per thread, or one polling thread that hands work to a pool (and then commits carefully).
  • Not closing the consumer. Without close(), the group waits for the session timeout before reassigning its partitions.
  • Using assign() and expecting rebalancing. Manually assigned consumers never rebalance; if one dies, its partitions are not read until you restart it.

In interviews

Know subscribe versus assign, that the consumer is single-threaded, and that committed offsets are “next to read”. Being able to sketch the poll loop from memory is often expected.

Poll loop

What it is. Everything a consumer does happens inside repeated poll() calls: fetching, delivering records to your code, and (with the classic protocol) parts of group membership. The loop is:

subscribe
loop:
    records = poll(timeout)
    for record in records: process(record)
    commit (if manual)
on shutdown: wakeup() from another thread -> poll throws WakeupException -> commit -> close()

How it works

  • poll() returns whatever is already buffered, or waits up to the timeout for a fetch to return. The consumer prefetches in the background, so the next poll often returns immediately.
  • max.poll.records (default 500) caps how many records one poll() returns. It does not change how much is fetched over the network.
  • max.poll.interval.ms (default 300000, 5 minutes) is the longest allowed gap between polls. If processing one batch takes longer, the consumer is considered stuck: it leaves the group, its partitions move to others, and its next commit fails.
  • Heartbeats are sent by a background thread, so long processing does not stop heartbeats. That is why there are two timeouts: the session timeout detects a dead process, and max.poll.interval.ms detects a live but stuck one.

Pitfalls

  • Slow processing causes rebalance loops. A consumer that takes 6 minutes per batch is evicted, rejoins, gets the same batch (its commit failed), and is evicted again. Lower max.poll.records, speed up processing, or raise max.poll.interval.ms deliberately.
  • Blocking calls inside the loop, such as one database round trip per record, are the usual cause of lag. Batch writes to the sink.
  • Handing records to other threads and polling on breaks commit ordering: you may commit offsets for records that are still being processed. Track per-partition completion, or use pause() while a batch is in progress.

In interviews

“Why does my consumer keep rebalancing?” A strong answer checks processing time against max.poll.interval.ms, and the session timeout and heartbeat settings, and mentions static membership and the new consumer protocol (see rebalancing).

Offset management

What it is. A consumer group’s progress is a set of committed offsets, one per partition, stored in the internal compacted topic __consumer_offsets. Separately, each consumer has a position: the next offset it will fetch, held in memory.

Term Meaning
Position next offset this consumer instance will read (in memory)
Committed offset next offset the group will read after a restart or rebalance (stored in Kafka)
Log-end offset (LEO) offset the next produced record will get
High watermark last offset replicated to all in-sync replicas, plus one; consumers can read only below it
Log start offset earliest offset still on disk (moves forward with retention)

What happens with no committed offset

When a group has no committed offset for a partition, or the committed offset has been deleted by retention, auto.offset.reset decides:

Value Behaviour
latest (default) start at the end: only new records
earliest start at the log start offset: everything still retained
by_duration:<ISO-8601 duration> start at the first record newer than now minus the duration, for example by_duration:PT6H (added in Kafka 4.0)
none throw an exception and let the application decide

Committed offsets also expire: when a group has no members for offsets.retention.minutes (default 10080, 7 days), its offsets are deleted, and the next start falls back to auto.offset.reset.

Pitfalls

  • The latest default skips history for a new group. Most data pipelines want earliest for their first run.
  • Committed offset versus processed offset. Commit last processed + 1. Committing the last processed offset itself reprocesses one record per partition on every restart.
  • Group offsets expire for groups that run rarely (a weekly batch job), silently resetting them.

In interviews

Explain where offsets live, the difference between position and committed offset, and auto.offset.reset. Mentioning that consumers only see records below the high watermark shows you understand the link to replication.

Auto vs manual commit

What it is. With enable.auto.commit=true (the default), the consumer commits offsets periodically (auto.commit.interval.ms, default 5 seconds). With false, your code calls commitSync() or commitAsync().

How auto-commit really behaves

In the Java client, auto-commit happens inside poll() (and close()): each poll may commit the positions reached by the previous poll. If you process every record fully before calling poll() again, auto-commit gives you at-least-once: a crash re-delivers up to five seconds of records. It becomes at-most-once (loss) if you hand records to another thread and poll again before they are processed.

librdkafka-based clients (including Confluent’s Python client) auto-commit from a background thread, and by default mark a record’s offset as ready to commit when it is handed to your application (enable.auto.offset.store=true). A crash during processing can therefore lose records. The common fix is enable.auto.offset.store=false and calling store_offsets(msg) after processing, keeping the background auto-commit.

Manual commit options

Call Behaviour Use
commitSync() blocks, retries until success or a non-retriable error after each batch; on shutdown and in onPartitionsRevoked
commitAsync(callback) returns at once, does not retry (a retry could overwrite a newer commit) in the hot loop, with commitSync() on shutdown
commitSync(Map<TopicPartition, OffsetAndMetadata>) commit specific offsets, for example per record or per partition fine-grained control

Worked example: commit before versus after processing

The simulation processes ten records in polls of four and crashes once, after six records have been written to the sink. It then restarts from the committed offset.

partition = [f"payment-{i}" for i in range(10)]   # offsets 0..9

def run(strategy, crash_after_processing, batch=4):
    """Process the partition in polls of `batch` records; crash once after processing
    `crash_after_processing` records; restart from the committed offset."""
    committed, sink, crashed, position = 0, [], False, 0
    while committed < len(partition):
        position = committed                      # (re)start from the committed offset
        while position < len(partition):
            records = partition[position:position + batch]
            if strategy == "commit-before":       # at-most-once
                committed = position + len(records)
            for r in records:
                if not crashed and len(sink) == crash_after_processing:
                    crashed = True
                    break                         # process dies here
                sink.append(r)
            else:
                position += len(records)
                if strategy == "commit-after":    # at-least-once
                    committed = position
                continue
            break                                 # crash: leave the inner loop and restart
    return sink

for strategy in ("commit-before", "commit-after"):
    sink = run(strategy, crash_after_processing=6)
    missing = sorted(set(partition) - set(sink), key=partition.index)
    dupes = sorted({r for r in sink if sink.count(r) > 1}, key=partition.index)
    print(f"{strategy:13}: wrote {len(sink):2} rows, missing {missing}, duplicated {dupes}")
commit-before: wrote  8 rows, missing ['payment-6', 'payment-7'], duplicated []
commit-after : wrote 12 rows, missing [], duplicated ['payment-4', 'payment-5']

Committing first lost the rest of the batch in flight. Committing after processing re-delivered the part of the batch that had been written before the crash. Neither is exactly-once; the at-least-once version is the one you can fix, by making the sink idempotent (upsert by payment_id) or by using transactions, as covered in delivery semantics.

Pitfalls

  • Committing per record with commitSync() is correct but slow: one round trip per record.
  • Retrying commitAsync() yourself can commit an older offset after a newer one.
  • Committing in the wrong rebalance callback: commit in onPartitionsRevoked, before the partitions move, or the new owner reprocesses them.

In interviews

“Is auto-commit safe?” Answer: in the Java client it is at-least-once if you finish processing before the next poll, and unsafe if processing is asynchronous; explain the Python client’s offset-store behaviour, and say you usually disable it for pipelines and commit after the sink write succeeds.

Consumer lag

What it is. Lag for a partition is the log-end offset minus the group’s committed offset: how many records the group has not yet processed. It is the main health signal for any streaming pipeline.

Worked example

After the Java consumer above stopped at 20 records:

bin/kafka-consumer-groups.sh --bootstrap-server localhost:19092 --describe --group payments-loader
bin/kafka-consumer-groups.sh --bootstrap-server localhost:19092 --describe --group payments-loader --state
GROUP           TOPIC           PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG             CONSUMER-ID     HOST            CLIENT-ID
payments-loader payments        0          2               12              10              -               -               -
payments-loader payments        1          4               4               0               -               -               -
payments-loader payments        2          14              14              0               -               -               -

GROUP           COORDINATOR (ID)          ASSIGNMENT-STRATEGY  STATE                #MEMBERS
payments-loader localhost:39092  (3)      uniform              Empty                0

The group is Empty (no running members), so there is no consumer ID. Partition 0 has 10 unprocessed records. uniform is the server-side assignor of the consumer protocol, because the consumer set group.protocol=consumer.

Measuring lag in production

Source What you get
kafka-consumer-groups.sh --describe point-in-time lag per partition
Consumer metrics records-lag, records-lag-max lag seen by the consumer itself
External exporters (for example Burrow, kafka-lag-exporter, or your platform’s monitoring) lag over time, alerts, lag in seconds

Lag in records is hard to alert on: 10,000 records may be one second on one topic and an hour on another. Alert on lag in time (age of the oldest unprocessed record) where your tooling supports it, and on lag that keeps growing.

Pitfalls

  • Lag stuck at a constant value with no members means the consumer is down.
  • Lag on one partition only usually means a hot key or a poison record blocking that partition.
  • Lag approaching retention means data will be deleted before it is read.

In interviews

“Your consumer lag keeps growing. What do you do?” Check whether consumers are running and not rebalancing, whether one partition is hot, where processing time goes (usually the sink), then scale consumers up to the partition count, batch sink writes, or add partitions.

Seek and rewind

What it is. Because Kafka keeps data, a consumer can move its position: back to replay after a bug fix, forward to skip bad data, or to a timestamp.

Worked example: seeking in code

Six readings were written to sensor-readings with explicit timestamps ten minutes apart from 00:00. The consumer then used assign() (no group) and seeked around:

c.assign(List.of(tp));                       // manual assignment: no group needed

c.seekToBeginning(List.of(tp));
System.out.println("seekToBeginning -> position " + c.position(tp));

c.seek(tp, 4);
System.out.println("seek(4)         -> first value " + c.poll(Duration.ofSeconds(2)).iterator().next().value());

long target = Instant.parse("2026-10-06T00:25:00Z").toEpochMilli();
OffsetAndTimestamp found = c.offsetsForTimes(Map.of(tp, target)).get(tp);
System.out.println("offsetsForTimes(00:25) -> offset " + found.offset()
        + " (record time " + Instant.ofEpochMilli(found.timestamp()) + ")");
c.seek(tp, found.offset());
for (ConsumerRecord<String, String> r : c.poll(Duration.ofSeconds(2)))
    System.out.println("  replay offset " + r.offset() + " " + r.value());

c.seekToEnd(List.of(tp));
System.out.println("seekToEnd       -> position " + c.position(tp));
seekToBeginning -> position 0
seek(4)         -> first value reading-4
offsetsForTimes(00:25) -> offset 3 (record time 2026-10-06T00:30:00Z)
  replay offset 3 reading-3
  replay offset 4 reading-4
  replay offset 5 reading-5
seekToEnd       -> position 6

offsetsForTimes returns the first offset whose timestamp is at or after the target. There was no record at 00:25, so it returned the 00:30 record.

Resetting a group’s offsets from the CLI

To replay for a whole group, stop its consumers first; the tool only resets inactive groups. Always --dry-run before --execute:

bin/kafka-consumer-groups.sh --bootstrap-server localhost:19092 --group payments-loader \
  --reset-offsets --to-earliest --topic payments --dry-run
bin/kafka-consumer-groups.sh --bootstrap-server localhost:19092 --group payments-loader \
  --reset-offsets --shift-by -2 --topic payments:2 --execute
GROUP           TOPIC           PARTITION  NEW-OFFSET
payments-loader payments        0          0
payments-loader payments        2          0
payments-loader payments        1          0

GROUP           TOPIC           PARTITION  NEW-OFFSET
payments-loader payments        2          12

Other targets include --to-offset, --to-datetime 2026-10-06T00:00:00.000, --by-duration PT1H and --to-latest. --topic payments:2 limits the reset to partition 2.

Pitfalls

  • Replays cause duplicates downstream unless the sink is idempotent or you truncate and reload the affected range.
  • Seeking before the log start offset triggers auto.offset.reset.
  • Seeking inside subscribe() before the assignment exists fails; seek in onPartitionsAssigned, or use assign().
  • Timestamps are producer-set by default (CreateTime), so a misbehaving producer can make time-based seeks land in the wrong place.

In interviews

“A bug corrupted yesterday’s output. How do you reprocess?” Fix and deploy the code, stop the group, reset offsets to a datetime with a dry run, make sure the sink overwrites rather than appends (or delete the bad range), restart, and watch lag fall. Mention that this only works within retention.

Deserialization

What it is. Kafka stores bytes. The consumer’s deserialisers (key.deserializer, value.deserializer) turn them back into objects. They must match the producer’s serialisers: String, Long, JSON, Avro or Protobuf with a Schema Registry, and so on.

Worked example: a poison pill

A bad producer wrote the string oops between two 8-byte longs on meter-totals. A consumer using LongDeserializer catches the error, which tells it the exact partition and offset, and skips that record:

try {
    for (ConsumerRecord<String, Long> r : c.poll(Duration.ofSeconds(1)))
        System.out.println("offset " + r.offset() + " value " + r.value());
} catch (RecordDeserializationException e) {
    System.out.println("cannot deserialise offset " + e.offset() + " in " + e.topicPartition()
            + ": " + e.getCause().getMessage());
    c.seek(e.topicPartition(), e.offset() + 1);   // skip it (better: send it to a dead-letter topic)
}
offset 0 value 100
cannot deserialise offset 1 in meter-totals-0: Size of data received by LongDeserializer is not 8
offset 2 value 250

Without the catch and seek, every poll() throws again at offset 1 and the partition is stuck forever.

Patterns for bad records

  • Dead-letter topic. Consume as bytes (ByteArrayDeserializer), deserialise in your own code, and send failures with error headers to a <topic>.dlq topic.
  • Schema Registry. Prevents most poison pills by rejecting incompatible schemas at produce time (see Schema Registry).
  • Framework support. Kafka Connect’s errors.tolerance and DLQ, and Kafka Streams’ deserialisation exception handlers, implement these patterns for you.

Pitfalls

  • Logging and continuing without moving the position: the same record fails forever.
  • Silently dropping bad records with no metric or DLQ: data loss nobody notices.
  • Different teams using different JSON conventions for the same topic. Agree a schema.

In interviews

“How do you handle a poison pill?” Explain that the consumer cannot move past it on its own, then describe catching the deserialisation error, writing the raw bytes to a DLQ with context, seeking past it, and alerting.

fetch.min.bytes and fetch.max.bytes

What it is. These settings control how much data a broker returns per fetch request, trading latency against efficiency.

Setting Default Meaning
fetch.min.bytes 1 broker waits until at least this much data is available…
fetch.max.wait.ms 500 …or this much time has passed, then replies anyway
fetch.max.bytes 52428800 (50 MiB) soft cap on one fetch response across all partitions
max.partition.fetch.bytes 1048576 (1 MiB) soft cap per partition per fetch
max.poll.records 500 records handed to your code per poll() (not a fetch setting)

The caps are soft: if the first record batch is bigger than the limit, it is still returned so the consumer can make progress.

Worked example: catching up on 20,000 records

A consumer read the 2 MB compress-none topic (20,000 records written in roughly 64 KiB batches) from the beginning with three settings, then reported its fetch-total and fetch-size-avg metrics:

fetch.min.bytes=1      max.partition.fetch.bytes=1048576  ->    2 fetches, avg  1018456 bytes, 20000 records
fetch.min.bytes=65536  max.partition.fetch.bytes=1048576  ->    2 fetches, avg  1018456 bytes, 20000 records
fetch.min.bytes=1      max.partition.fetch.bytes=102400   ->   31 fetches, avg    65707 bytes, 20000 records

When a consumer is catching up, data is always available, so fetch.min.bytes makes no difference and the per-partition cap decides how many round trips are needed. When a consumer is tailing a quiet topic, fetch.min.bytes matters: with the default of 1 the broker replies as soon as any record arrives, and a larger value makes it wait (up to fetch.max.wait.ms) for a fuller response, which cuts requests and broker CPU but adds latency.

Tuning guidance

Goal Change
Lower broker load from many idle consumers raise fetch.min.bytes (for example tens of KB)
Lowest latency keep fetch.min.bytes=1, lower fetch.max.wait.ms
Faster catch-up on many partitions raise max.partition.fetch.bytes, keep fetch.max.bytes high enough
Consumer running out of memory lower max.partition.fetch.bytes (memory scales with partitions assigned times this)

Pitfalls

  • Large records must fit within the broker and topic limits; consumer fetch limits are soft, so they rarely block progress today, but memory use grows with them.
  • Raising max.poll.records to “fetch more” changes nothing on the network; it only changes how many records each poll hands to your code.

In interviews

Explain the pair fetch.min.bytes / fetch.max.wait.ms as a latency-versus-efficiency knob, max.partition.fetch.bytes as the per-partition cap that also bounds memory, and that max.poll.records is a processing knob, not a fetch knob.

Practice questions

Your consumer commits offsets with auto-commit and hands records to a thread pool. After a crash, some records were never processed. Why?

Auto-commit commits the positions reached by earlier polls. Because processing happened on other threads, the main loop kept polling, so offsets for records still in the pool were committed. When the process crashed, those records were lost. Disable auto-commit and commit only offsets whose records have finished processing, per partition, or process synchronously in the poll loop.

A new consumer group starts and reads nothing, although the topic has a week of data. What is wrong?

The group has no committed offsets, so auto.offset.reset applies, and its default is latest: start at the end. Set auto.offset.reset=earliest (or by_duration:... on Kafka 4.0+ clients), or reset the group’s offsets with kafka-consumer-groups.sh --reset-offsets --to-earliest before starting.

One partition’s lag keeps growing while others are at zero. What would you check?

A hot key concentrating traffic in that partition, a poison record that the consumer keeps failing on, or slow processing for that partition’s data (for example large records or a slow downstream key). Check producer key distribution, consumer logs for repeated errors at the same offset, and processing time per partition.

What is the difference between max.poll.interval.ms and session.timeout.ms?

The session timeout detects a dead consumer process: heartbeats (sent from a background thread) stop arriving. max.poll.interval.ms detects a live but stuck consumer: it is still heartbeating but has not called poll() within the limit, usually because processing a batch takes too long. Either one removes the consumer from the group and triggers a rebalance. With the consumer protocol, the session timeout is a broker setting (group.consumer.session.timeout.ms).

How do you reprocess the last six hours of a topic for one consumer group?

Stop the group’s consumers, run kafka-consumer-groups.sh --reset-offsets --by-duration PT6H --group ... --topic ... --dry-run, check the new offsets, rerun with --execute, and restart. Make sure the sink is idempotent or that the affected output is removed first, otherwise the replay creates duplicates.

Why should you not retry a failed commitAsync() call?

Asynchronous commits can complete out of order. If an older commit fails and you retry it after a newer commit has succeeded, the retry moves the committed offset backwards, causing reprocessing after the next restart. Let later commits supersede it, and use commitSync() on shutdown and during rebalances.

Key takeaways

  • The consumer is a single-threaded poll loop; poll() fetches, delivers records and keeps the group membership alive.
  • Committed offsets live in __consumer_offsets, mean “next to read”, and fall back to auto.offset.reset (default latest) when missing or expired.
  • Commit after processing for at-least-once, and make the sink idempotent; committing before processing risks loss.
  • Lag is log-end offset minus committed offset; alert on growth and on lag in time.
  • Replay with seek, offsetsForTimes or kafka-consumer-groups.sh --reset-offsets, within retention, with a dry run first.
  • A record that cannot be deserialised blocks its partition until you catch the error, store it in a DLQ and seek past it.
  • fetch.min.bytes and fetch.max.wait.ms trade latency for efficiency; max.partition.fetch.bytes caps per-partition fetches and memory.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Java consumer examples were run against a three-broker Apache Kafka 4.3.1 KRaft cluster with kafka-clients 4.3.1 on Java 21; the Python consumer was run with confluent-kafka 2.15.1. Those blocks are marked noexec because they need a broker. kafka-consumer-groups.sh output is from the same cluster. The commit-order simulation runs on plain Python 3.

Progress is saved in this browser only. No account needed.

Search
Filter by type