Learn Labs
4. Kafka Consumers: Reading Data from Kafka

4.7 Commits and offsets — where consumers actually break

Get this wrong by one and you reprocess one record per commit forever (harmless-ish) or skip one record per commit (data loss).

Commit ordering

commit
3
duplicates
0
lost
poll
5 records
process
3 donecrash
commit
never reached
Safe

At-least-once. The offset was never committed, so after the rebalance another consumer reprocesses all 5 records — 3 of which were already handled. This is the right default: make the processing idempotent and duplicates stop mattering.

A consumer crash between processing and committing is not an edge case — it is the normal case during a deploy. The order of the two decides which failure you get, and you must pick one.

7.1 The model

"One of Kafka's unique characteristics is that it does not track acknowledgments from consumers the way many JMS queues do. Instead, it allows consumers to use Kafka to track their position (offset) in each partition."

"Unlike traditional message queues, Kafka does not commit records individually. Instead, consumers commit the last message they've successfully processed from a partition and implicitly assume that every message before the last was also successfully processed."

Mechanically: the consumer "sends a message to Kafka, which updates a special __consumer_offsets topic with the committed offset for each partition."

commit offset 5000 for T-p3Consumer__consumer_offsetsa regular topic!Kafka
Figure 4.7.17.1 The model

Why any of this matters: "As long as all your consumers are up, running, and churning away, this will have no impact. However, if a consumer crashes or a new consumer joins the consumer group, this will trigger a rebalance. After a rebalance, each consumer may be assigned a new set of partitions. In order to know where to pick up the work, the consumer will read the latest committed offset of each partition and continue from there."

7.2 The two failure modes — this is the whole ballgame

Case A — committed offset < last processed offset
log ……100101102103104105106107committed = 102processed through 106THESE ARE PROCESSED TWICE ► DUPLICATES
Case B — committed offset > last processed offset
log ……100101102103104105106107processed through 102committed = 106THESE ARE NEVER PROCESSED ► MISSED MESSAGES

Silent data loss.

Figure 4.7.27.2 The two failure modes — this is the whole ballgame

"Clearly, managing offsets has a big impact on the client application."

7.3 ⚠️ "Which offset is committed?" — the off-by-one everyone hits

"When committing offsets either automatically or without specifying the intended offsets, the default behavior is to commit the offset AFTER the last offset that was returned by poll()."

The book's convention: "it is tedious to repeatedly read 'Commit the offset that is one larger than the last offset the client received from poll(),' and 99% of the time it does not matter. So, we are going to write 'Commit the last offset' when we refer to the default behavior — and if you need to manually manipulate offsets, please keep this note in mind."

The rule when committing manually:

 committed offset = offset of the NEXT message your application will read
                  = last_processed_offset + 1

Get this wrong by one and you reprocess one record per commit forever (harmless-ish) or skip one record per commit (data loss).

7.4 Automatic commit

enable.auto.commit = true
auto.commit.interval.ms = 5000   // default: 5 seconds

Mechanics: "Just like everything else in the consumer, the automatic commits are driven by the poll loop. Whenever you poll, the consumer checks if it is time to commit, and if it is, it will commit the offsets it returned in the last poll."

The duplicate window:

3 SECONDS OF EVENTS PROCESSED TWICEcommitoffset 1000processprocess✗ CONSUMER CRASHEShaving processed through offset 1750

Rebalance; survivors start from the LAST COMMITTED offset = 1000.

Figure 4.7.3The duplicate window

"It is possible to configure the commit interval to commit more frequently and reduce the window in which records will be duplicated, but it is impossible to completely eliminate them."

⚠️ The subtle killer:

"With autocommit enabled, when it is time to commit offsets, the next poll will commit the last offset returned by the previous poll. It doesn't know which events were actually processed, so it is critical to always process all the events returned by poll() before calling poll() again. (Just like poll(), close() also commits offsets automatically.) This is usually not an issue, but pay attention when you handle exceptions or exit the poll loop prematurely."

records = poll()           // returns offsets 1000..1500
for (r : records) {
    process(r);            // ✗ throws at offset 1200
}                          //   you `catch` and `continue` the while loop
// next poll() → COMMITS 1500.  Offsets 1201..1500 were NEVER PROCESSED.
//                              SILENT DATA LOSS.

Verdict: "Automatic commits are convenient, but they don't give developers enough control to avoid duplicate messages."

7.5 commitSync() — commit current offset

Set enable.auto.commit=false. "The simplest and most reliable of the commit APIs is commitSync(). This API will commit the latest offset returned by poll() and return once the offset is committed, throwing an exception if the commit fails."

Duration timeout = Duration.ofMillis(100);

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(timeout);
    for (ConsumerRecord<String, String> record : records) {
        System.out.printf("topic = %s, partition = %d, offset = %d, " +
            "customer = %s, country = %s\n",
            record.topic(), record.partition(),
            record.offset(), record.key(), record.value());
    }
    try {
        consumer.commitSync();      // AFTER processing the whole batch
    } catch (CommitFailedException e) {
        log.error("commit failed", e);
    }
}

Positioning is everything:

"commitSync() will commit the latest offset returned by poll(), so if you call commitSync() before you are done processing all the records in the collection, you risk missing the messages that were committed but not processed, in case the application crashes. If the application crashes while it is still processing records in the collection, all the messages from the beginning of the most recent batch until the time of the rebalance will be processed twice — this may or may not be preferable to missing messages."

Retry behavior: "commitSync retries committing as long as there is no error that can't be recovered. If this happens, there is not much we can do except log an error."

Note also: "You should determine when you are 'done' with a record according to your use case." — the book is explicit that "processed" is your definition, not Kafka's.

7.6 commitAsync()

The problem with sync: "the application is blocked until the broker responds to the commit request. This will limit the throughput of the application. Throughput can be improved by committing less frequently, but then we are increasing the number of potential duplicates that a rebalance may create."

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(timeout);
    for (ConsumerRecord<String, String> record : records) { /* process */ }
    consumer.commitAsync();      // send and carry on
}
⚠️ Why commitAsync() deliberately does NOT retry

This is one of the most instructive explanations in the book:

consumerbrokercommit offset 2000t0 — lostcommit offset 3000t2 — succeedsretry of 2000t3 — too late
  • t0: a temporary communication problem means the broker NEVER GETS the request, NEVER RESPONDS. t1: the consumer processes another batch. t2: offset 3000 commits successfully — the current position is 3000.
  • t3: if commitAsync() RETRIED the failed 2000 commit, it might succeed after 3000 was already processed and committed, so the COMMITTED OFFSET GOES BACKWARDS (3000 → 2000).
  • “In the case of a rebalance, this will cause MORE duplicates.”
Figure 4.7.5This is one of the most instructive explanations in the book

"The reason it does not retry is that by the time commitAsync() receives a response from the server, there may have been a later commit that was already successful."

With a callback:

consumer.commitAsync(new OffsetCommitCallback() {
    public void onComplete(Map<TopicPartition, OffsetAndMetadata> offsets,
                           Exception e) {
        if (e != null)
            log.error("Commit failed for offsets {}", offsets, e);
    }
});

"It is common to use the callback to log commit errors or to count them in a metric, but if you want to use the callback for retries, you need to be aware of the problem with commit order."

Retrying async commits safely — the sequence-number pattern

"A simple pattern to get the commit order right for asynchronous retries is to use a monotonically increasing sequence number. Increase the sequence number every time you commit, and add the sequence number at the time of the commit to the commitAsync callback. When you're getting ready to send a retry, check if the commit sequence number the callback got is equal to the instance variable; if it is, there was no newer commit and it is safe to retry. If the instance sequence number is higher, don't retry because a newer commit was already sent."

7.7 Combining sync and async — the standard production shape

The insight: "Normally, occasional failures to commit without retrying are not a huge problem, because if the problem is temporary, the following commit will be successful. But if we know that this is the LAST commit before we close the consumer, or before a rebalance, we want to make extra sure that the commit succeeds."

Duration timeout = Duration.ofMillis(100);

try {
    while (!closing) {
        ConsumerRecords<String, String> records = consumer.poll(timeout);
        for (ConsumerRecord<String, String> record : records) {
            /* process */
        }
        consumer.commitAsync();      // ① fast; a failure is retried by the
                                     //    NEXT commit
    }
    consumer.commitSync();           // ② no "next commit" exists — retry
                                     //    until success or unrecoverable
} catch (Exception e) {
    log.error("Unexpected error", e);
} finally {
    consumer.close();
}
WhenUseWhy
STEADY STATEcommitAsync()fast; the next commit is the retry
SHUTDOWNcommitSync()no next commit; it must land
BEFORE REBALANCEcommitSync() in onPartitionsRevoked()the handover has to see the offsets

7.8 Committing a specified offset

Why you'd need it: "Committing the latest offset only allows you to commit as often as you finish processing batches. But what if you want to commit more frequently? What if poll() returns a huge batch and you want to commit offsets in the middle of the batch to avoid having to process all those rows again if a rebalance occurs? You can't just call commitSync() or commitAsync() — this will commit the last offset returned, which you didn't get to process yet."

private Map<TopicPartition, OffsetAndMetadata> currentOffsets = new HashMap<>();
int count = 0;
Duration timeout = Duration.ofMillis(100);

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(timeout);
    for (ConsumerRecord<String, String> record : records) {
        System.out.printf("topic = %s, partition = %s, offset = %d, " +
            "customer = %s, country = %s\n",
            record.topic(), record.partition(), record.offset(),
            record.key(), record.value());

        currentOffsets.put(
            new TopicPartition(record.topic(), record.partition()),
            new OffsetAndMetadata(record.offset()+1, "no metadata"));
            //                                  ▲▲
            //   THE COMMITTED OFFSET SHOULD ALWAYS BE THE OFFSET OF
            //   THE NEXT MESSAGE YOUR APPLICATION WILL READ

        if (count % 1000 == 0)
            consumer.commitAsync(currentOffsets, null);   // no callback → null
        count++;
    }
}

Cost: "Since your consumer may be consuming more than a single partition, you will need to track offsets on all of them, which adds complexity to your code."

You can commit "based on time or perhaps content of the records" — the count % 1000 is just one policy. And commitSync() is "also completely valid here."


On this page