Kafka courseLesson 4 of 6
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.
On this page
- Sample data
- Consumer API
- Worked example: Java consumer with manual commits
- Worked example: Python consumer
- Pitfalls
- In interviews
- Poll loop
- How it works
- Pitfalls
- In interviews
- Offset management
- What happens with no committed offset
- Pitfalls
- In interviews
- Auto vs manual commit
- How auto-commit really behaves
- Manual commit options
- Worked example: commit before versus after processing
- Pitfalls
- In interviews
- Consumer lag
- Worked example
- Measuring lag in production
- Pitfalls
- In interviews
- Seek and rewind
- Worked example: seeking in code
- Resetting a group’s offsets from the CLI
- Pitfalls
- In interviews
- Deserialization
- Worked example: a poison pill
- Patterns for bad records
- Pitfalls
- In interviews
- fetch.min.bytes and fetch.max.bytes
- Worked example: catching up on 20,000 records
- Tuning guidance
- Pitfalls
- In interviews
- Practice questions
- 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
ConcurrentModificationExceptionin 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 onepoll()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.msdetects 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 raisemax.poll.interval.msdeliberately. - 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
latestdefault skips history for a new group. Most data pipelines wantearliestfor 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 inonPartitionsAssigned, or useassign(). - 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>.dlqtopic. - Schema Registry. Prevents most poison pills by rejecting incompatible schemas at produce time (see Schema Registry).
- Framework support. Kafka Connect’s
errors.toleranceand 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.recordsto “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 toauto.offset.reset(defaultlatest) 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,offsetsForTimesorkafka-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.bytesandfetch.max.wait.mstrade latency for efficiency;max.partition.fetch.bytescaps per-partition fetches and memory.
Progress is saved in this browser only. No account needed.