A Kafka consumer looks simple: subscribe to a topic, call poll() in a loop, process records. Most production incidents with Kafka consumers come from what that loop hides: a background fetcher, a liveness deadline, an offset that is committed at a moment you did not choose, and a group protocol that moves partitions between instances while your code is mid-batch. Understanding those pieces is the difference between a pipeline that recovers by itself and one that duplicates, skips or stalls.
The log model, partitions and producer guarantees are introduced in the Apache Kafka overview. This article goes inside the consumer: the poll loop, fetching, commits, both group protocols, threading and lag. Configuration names and defaults below were checked against the Apache Kafka 4.x documentation.
The model: positions, commits and a coordinator
A partition is an append-only log, and a consumer's progress in it is just a number: the offset of the next record it will read. Kafka tracks two such numbers per partition. The position lives in the consumer's memory and advances as poll() returns records. The committed offset is stored durably by the cluster in the internal __consumer_offsets topic, keyed by group, topic and partition, and is where a consumer starts after a restart or a reassignment. Everything about delivery guarantees is about when the second number catches up with the first.
A consumer group is a set of consumers sharing a group.id. Kafka assigns each partition of the subscribed topics to exactly one member, so the group divides the work and each partition is processed in order by one consumer at a time. One broker acts as the group coordinator for each group: it tracks membership, drives assignment and stores commits. Different groups reading the same topic are completely independent, each with its own offsets, which is how one topic feeds a search indexer, a fraud model and an archive at once.
What poll() actually does
In the Java client, poll(timeout) returns records already fetched in the background where possible, and sends fetch requests to the partition leaders for more. A fetch asks the broker to wait until at least fetch.min.bytes are available (default 1 byte) or fetch.max.wait.ms has passed (default 500 ms), and returns at most max.partition.fetch.bytes per partition (default 1 MiB). poll() then hands your code at most max.poll.records records (default 500).
Two timers govern liveness. The consumer heartbeats to the coordinator from a background thread; if no heartbeat arrives within the session timeout, the member is considered dead. Separately, if your code does not call poll() again within max.poll.interval.ms (default 300,000 ms, five minutes), the client itself leaves the group, because a consumer that heartbeats but never polls is alive but stuck. The first timer catches crashed processes; the second catches slow or hung processing. The second one is the one that bites, as the worked example shows.
The consumer object is not thread-safe. Only wakeup() may be called from another thread, which makes a blocked poll() throw WakeupException so the loop can exit and close cleanly.
Offset commits and delivery guarantees
With the default enable.auto.commit=true, the client commits the offsets of records returned by earlier polls every auto.commit.interval.ms (default 5,000 ms), as part of later poll() calls and on close. If you process every record synchronously before polling again, that gives at-least-once delivery: a crash replays up to a few seconds of records. If you hand records to another thread and poll on, auto-commit can commit records that were never processed, and a crash loses them. That is the most common way teams lose data with Kafka.
Manual commits make the moment explicit. Commit after processing for at-least-once; the committed value is the offset of the next record to read, so last processed offset plus one. commitSync blocks and retries; commitAsync does not block and does not retry, so the usual pattern is async commits in the loop and a sync commit on shutdown and when partitions are revoked. Exactly-once needs more than commit ordering: the output and the offset must commit atomically, either with Kafka transactions (sendOffsetsToTransaction and consumers using isolation.level=read_committed) or by storing offsets in the same database transaction as the results. Both are covered in exactly-once semantics.
Properties props = new Properties();
props.put("bootstrap.servers", "kafka:9092");
props.put("group.id", "order-enricher");
props.put("enable.auto.commit", "false");
props.put("auto.offset.reset", "earliest"); // only used when the group has no committed offset
props.put("max.poll.records", "200");
props.put("group.instance.id", System.getenv("POD_NAME")); // static membership
props.put("key.deserializer", StringDeserializer.class.getName());
props.put("value.deserializer", StringDeserializer.class.getName());
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
Map<TopicPartition, OffsetAndMetadata> done = new HashMap<>();
consumer.subscribe(List.of("orders"), new ConsumerRebalanceListener() {
public void onPartitionsRevoked(Collection<TopicPartition> parts) {
consumer.commitSync(done); // flush progress before losing partitions
parts.forEach(done::remove);
}
public void onPartitionsAssigned(Collection<TopicPartition> parts) { }
});
try {
while (running) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> r : records) {
enrichAndWrite(r); // must be idempotent: replays happen
done.put(new TopicPartition(r.topic(), r.partition()),
new OffsetAndMetadata(r.offset() + 1));
}
consumer.commitAsync(new HashMap<>(done), null);
}
} catch (WakeupException e) {
// shutdown requested from another thread via consumer.wakeup()
} finally {
try { consumer.commitSync(done); } finally { consumer.close(); }
}Two settings decide what happens with no usable commit. auto.offset.reset (default latest) applies when a group has never committed for a partition or its commit is gone; latest silently skips everything already in the log, which surprises every new pipeline that expected a backfill. And committed offsets are deleted once a group has been empty for offsets.retention.minutes (broker default 10,080, seven days), so a consumer that was off for a long holiday resets too. Set the reset policy deliberately, or set it to none and handle the exception.
Two group protocols
Kafka now has two ways for a group to agree on an assignment, selected on the client by group.protocol. The default is still classic; the new consumer protocol from KIP-848 became generally available in Kafka 4.0 and is enabled on 4.0 brokers by default.
In the classic protocol, members send JoinGroup, the coordinator picks one member as leader, the leader computes the assignment client-side with the configured assignor, and SyncGroup distributes it. With an eager assignor every member gives up all partitions at the start of each rebalance, so the whole group pauses. Cooperative assignment revokes only the partitions that move, over two rounds. Note the default partition.assignment.strategy: it lists RangeAssignor and then CooperativeStickyAssignor, and the group uses the first assignor all members support, so out of the box you get eager range assignment. The cooperative entry is there to allow a rolling upgrade to it; you have to remove Range in a second rollout to finish the switch.
In the consumer protocol, members send a periodic ConsumerGroupHeartbeat, and the coordinator computes the assignment server-side and reconciles each member incrementally. There is no group-wide barrier, so one slow member no longer stalls everyone. Assignors are broker plugins, uniform (default) and range, chosen per member with group.remote.assignor; custom client-side assignors are not supported. The client settings session.timeout.ms, heartbeat.interval.ms and partition.assignment.strategy no longer apply.
| classic | consumer (KIP-848) | |
|---|---|---|
| Selected by | group.protocol=classic (default) | group.protocol=consumer |
| Assignment computed by | leader member, client-side | group coordinator, server-side |
| Rebalance style | eager or cooperative per assignor | incremental per member, no global barrier |
| Session timeout | client session.timeout.ms, default 45 s | broker group.consumer.session.timeout.ms, default 45 s, minimum 45 s |
| Heartbeat interval | client heartbeat.interval.ms, default 3 s | broker group.consumer.heartbeat.interval.ms, default 5 s |
| Poll deadline | max.poll.interval.ms, client | max.poll.interval.ms, client |
| Migration | online rolling upgrade from a classic group is supported |
One operational consequence: under the new protocol the broker's minimum session timeout is 45 seconds by default, so a team that tuned a 10-second session for fast failover on the classic protocol cannot carry that over without a broker change. Rebalance mechanics themselves, including cooperative rounds and their pitfalls, are in consumer group rebalancing.
Static membership
By default a restarted consumer is a new member, so a rolling deploy of ten pods causes ten departures and ten joins. Setting group.instance.id to a stable per-instance name, such as the pod name of a StatefulSet, makes the member static: when it disappears, the coordinator keeps its partitions reserved until the session timeout expires, and if it returns in time it gets the same partitions back with no rebalance. The trade-off is that a truly dead static member is noticed only after the session timeout, so its partitions sit unread for that long. Choose a timeout longer than a normal restart and shorter than the lag you can tolerate.
Threading and parallelism
Parallelism within a group is capped by partitions: a group with more members than partitions leaves the extras idle. Within that cap there are two designs. One consumer per thread is the simplest and keeps per-partition order trivially, at the cost of one set of connections and buffers per thread. One consumer feeding a worker pool raises concurrency beyond the partition count, but you then own ordering and commits: route records by key to a fixed worker if per-key order matters (see Kafka ordering), call pause() on partitions whose queues are full so the loop keeps polling and stays in the group, and commit only offsets below the lowest record still in flight.
class PartitionTracker:
"""Committable offset = lowest offset still in flight, or next offset if none."""
def __init__(self):
self.in_flight = {} # partition -> sorted list of offsets being processed
self.next_offset = {} # partition -> highest offset seen + 1
def started(self, part, off):
self.in_flight.setdefault(part, []).append(off)
self.next_offset[part] = max(self.next_offset.get(part, 0), off + 1)
def finished(self, part, off):
self.in_flight[part].remove(off)
def committable(self, part):
pending = self.in_flight.get(part)
return min(pending) if pending else self.next_offset.get(part)
Worked example: the five-minute cliff
An enricher group reads a 12-partition topic at 2,000 records per second and calls a pricing service for each record, taking about 8 ms. Throughput needed is 2,000 x 0.008 = 16 records in flight at any moment. With one consumer per partition and synchronous processing, 12 consumers can do 12 x 125 = 1,500 per second, so lag grows by 500 per second. Either raise partitions to at least 16 (more, for headroom) or add keyed concurrency inside each consumer.
Now the pricing service degrades to 800 ms per call. A poll returns up to 500 records, and 500 x 0.8 s = 400 s, which exceeds the 300 s poll deadline. The consumer leaves the group mid-batch, its partitions move to a peer that receives the same slow batch from the last commit, and the group enters a loop of rebalances while processing nothing new and replaying the same records. The fix is to bound batch time, not just raise the deadline: set max.poll.records so that records times worst-case latency stays well under max.poll.interval.ms, put a timeout on the downstream call, and use pause and resume when a dependency is slow.
Lag and monitoring
Consumer lag per partition is the log end offset minus the committed offset; with read_committed the relevant end is the last stable offset. The quickest view is kafka-consumer-groups.sh --bootstrap-server kafka:9092 --describe --group order-enricher, which lists current offset, log end offset, lag and owner for each partition. Alert on lag that keeps growing and on time lag (how old the oldest unprocessed record is), because 50,000 records of lag is seconds on one topic and hours on another. Also watch rebalance rate, commit failures and partitions with no owner. Offsets can be moved with --reset-offsets plus a target such as --to-datetime and --execute, only while the group has no active members.
Failure modes
- Auto-commit with asynchronous processing commits work that never finished; a crash loses it.
- Poll deadline exceeded by slow batches causes repeated rebalances and replays.
- auto.offset.reset=latest on a new group skips the backlog without any error.
- Offsets expired after a group was empty past retention, so it restarts from the reset policy.
- Hot partitions from skewed keys leave one consumer far behind while others idle; adding consumers does not help. See partitioning strategies.
- Non-idempotent side effects: every at-least-once pipeline replays after rebalances, so writes must be upserts or deduplicated.
- Mixed assignor lists during a rollout leave the group on the old eager protocol longer than intended.
What to do next
- List each consumer group and write down its commit mode, reset policy,
max.poll.recordsand the worst-case time to process one batch. - Turn off auto-commit wherever processing is asynchronous, and commit offset plus one after processing, with a sync commit on revoke and close.
- Set
auto.offset.resetdeliberately for every group, and alert when a group resets. - Add static membership to consumers deployed as rolling replicas, with a session timeout that covers a normal restart.
- Test
group.protocol=consumerin staging on 4.x brokers, removing client-side timeout and assignor settings, and compare rebalance duration during a rolling deploy. - Dashboard time lag, rebalance rate and unowned partitions per group, and make downstream writes idempotent.