Ch. 26 · Apache Kafka

Kafka Offset Commits and Delivery Semantics Explained

How Kafka consumers commit offsets, why commit order decides at-most-once or at-least-once, and how to get exactly-once effects.

~8 min readintermediateupdated Oct 6, 2026

“Explain at-most-once, at-least-once and exactly-once in Kafka” is one of the most predictable Kafka interview questions, and the strongest answers tie it to one concrete thing: when the consumer commits its offset relative to when it processes the record. Interviewers then probe the details that cause real incidents: “is auto-commit safe?”, “why does commitAsync not retry?”, “what happens to uncommitted work during a rebalance?” and “how do you avoid double-charging a customer?”

This guide explains what an offset commit is, how each commit style behaves under a crash, and how to build exactly-once effects on top of at-least-once delivery. It targets Kafka 3.8/4.0 and the Java client. The Java code is illustrative and was not run against a cluster; the crash simulation in the verification section is a small Node.js script that was run.

Before you start

You should know that a consumer group divides partitions among its members and that each partition is an ordered log addressed by offsets. If not, start with Kafka partitions, offsets and ordering. Comfort with Java and the idea of a database transaction will help with the later steps.

The short answer

A consumer group stores a committed offset per partition in the internal __consumer_offsets topic; it is the offset of the next record to read after a restart or rebalance. If you commit after processing, a crash between the two replays records, which is at-least-once. If you commit before processing, a crash skips records, which is at-most-once. Kafka’s exactly-once covers read-process-write pipelines that stay inside Kafka, using transactions. When the side effect is a database write or an API call, you get exactly-once effects by making that effect idempotent, or by storing the offset in the same database transaction.

How it works

A consumer has two separate notions of progress for each partition:

  • Position: the offset of the next record poll() will return. It lives in memory and advances as soon as records are returned to your code, whether or not you have processed them.
  • Committed offset: the position the group has saved to the broker. It survives restarts and is what a new owner of the partition starts from after a rebalance.

Commits go to the group coordinator, which appends them to __consumer_offsets, a compacted topic keyed by group, topic and partition. Because the committed offset is “where to resume”, after processing offset 41 you commit 42. The no-argument commitSync() and commitAsync() commit the current position for every assigned partition, which is the offset after the last record returned by poll().

There are three ways to commit:

Style Behaviour Typical risk
Auto-commit (enable.auto.commit=true, the default) during poll() and close(), at most every auto.commit.interval.ms (5 s), commits the positions from earlier polls loss if records are processed on another thread
commitSync() blocks until the commit succeeds, retrying retriable errors latency per commit
commitAsync(callback) sends the commit and returns; does not retry a failed commit is only logged

commitAsync does not retry by design. If a commit for offset 2000 fails and is retried after a commit for 3000 has succeeded, the retry would move the group backwards. Committed offsets also expire: once a group has no members, its offsets are deleted after offsets.retention.minutes (7 days), after which auto.offset.reset decides where it starts.

Step-by-step walkthrough

Step 1: Turn off auto-commit and commit after processing

Explicit commits make the guarantee visible in the code. Processing the whole batch and then committing gives at-least-once with one commit per poll.

props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);

while (running) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
    for (ConsumerRecord<String, String> r : records) {
        handle(r);                       // must be idempotent: it can run twice after a crash
    }
    if (!records.isEmpty()) consumer.commitAsync((offsets, ex) -> {
        if (ex != null) log.warn("commit failed for {}", offsets, ex);
    });
}
java

The window for duplicates is everything processed since the last successful commit, so a crash replays at most one batch per partition.

Step 2: Commit synchronously on shutdown and on revoke

An async commit in the loop is fast, but the last one may fail. Two moments deserve a blocking commit: when this instance shuts down, and when a rebalance is about to take partitions away.

Runtime.getRuntime().addShutdownHook(new Thread(consumer::wakeup));
try {
    while (true) { /* poll, handle, commitAsync as in Step 1 */ }
} catch (WakeupException e) {
    // expected on shutdown
} finally {
    try { consumer.commitSync(); } finally { consumer.close(); }
}
java

For revocation, call commitSync() inside ConsumerRebalanceListener.onPartitionsRevoked, so the next owner starts exactly where this one stopped. The shutdown hook only calls wakeup(), the one thread-safe consumer method, and the polling thread then exits the loop cleanly.

Step 3: Commit specific offsets when batches are large

If a batch takes a while, committing per partition as you go narrows the replay window. Remember the plus one.

Map<TopicPartition, OffsetAndMetadata> done = new HashMap<>();
int count = 0;
for (ConsumerRecord<String, String> r : records) {
    handle(r);
    done.put(new TopicPartition(r.topic(), r.partition()), new OffsetAndMetadata(r.offset() + 1));
    if (++count % 100 == 0) consumer.commitAsync(Map.copyOf(done), null);
}
consumer.commitSync(done);
java

Step 4: Make the side effect idempotent

At-least-once means a record can be handled twice, so the handler must make the second run harmless. The simplest form is a unique constraint on a stable event ID carried in the record.

INSERT INTO payments (event_id, order_id, amount_cents)
VALUES (:eventId, :orderId, :amount)
ON CONFLICT (event_id) DO NOTHING;
sql

A replayed record now inserts nothing. This turns at-least-once delivery into exactly-once effects for that table without any Kafka transaction.

Step 5: Or store offsets with the data

When the output lives in one database, you can make the database the source of truth for progress. Write the business rows and the partition’s next offset in one transaction, and on assignment seek to the stored offset instead of using Kafka’s committed one.

// Illustrative
@Override public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
    for (TopicPartition tp : partitions) {
        consumer.seek(tp, offsetStore.nextOffset(tp));    // read from the same database
    }
}

void handleInTransaction(ConsumerRecord<String, String> r) {
    db.inTransaction(tx -> {
        tx.apply(r.value());
        tx.saveOffset(r.topic(), r.partition(), r.offset() + 1);
    });
}
java

A crash either commits both the change and the offset or neither, so nothing is lost or duplicated.

Worked scenario

A notification service consumed order-shipped events with auto-commit on and, to go faster, handed each record to a thread pool that sent the email. Latency improved, and then a deploy killed a pod while 300 emails were queued in the pool. Those customers never got their notification.

// Broken: auto-commit + work handed to another thread
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true);
for (ConsumerRecord<String, String> r : consumer.poll(Duration.ofMillis(500))) {
    pool.submit(() -> emailService.send(r.value()));
}
java

Auto-commit commits the positions returned by previous polls. The poll loop raced ahead of the pool, so offsets for queued but unsent emails had already been committed. When the pod died, the group resumed after them: at-most-once by accident.

The fix keeps the parallelism but commits only what is finished. Track completed offsets per partition and commit the highest offset below which everything is done:

// Fixed: commit only the contiguous completed prefix per partition
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
for (ConsumerRecord<String, String> r : records) {
    tracker.started(r);
    pool.submit(() -> { emailService.send(r.value()); tracker.finished(r); });
}
consumer.commitAsync(tracker.safeOffsets(), null);   // e.g. 0-9 done, 10 pending, 11-20 done -> commit 10
java

A crash now replays some emails instead of dropping them, so the email service also deduplicates by event ID.

Common mistake

  • “At-least-once is the default, so I do not need to think about it.” It is only at-least-once if processing completes before the commit.
  • “auto.offset.reset=earliest makes a restarted consumer start from the beginning.” It applies only when there is no valid committed offset.
  • Committing record.offset() instead of record.offset() + 1. The last record is replayed after every restart.
  • “Kafka transactions make my database write exactly-once.” They cover writes to Kafka and offset commits, not external systems.
  • Retrying commitAsync manually in its callback without checking ordering. An old offset can overwrite a newer one.

Verify the behavior

This Node.js simulation reads offsets 0 to 9 in batches of 4 and crashes once, right after processing offset 5. It printed the output in the comments when run with Node 22:

function run(strategy) {
  const processed = [];
  let committed = 0, crashed = false;
  while (committed < 10) {
    const batch = [];
    for (let o = committed; o < Math.min(committed + 4, 10); o++) batch.push(o);
    const next = batch.at(-1) + 1;
    if (strategy === 'commit-first') committed = next;
    let failed = false;
    for (const o of batch) {
      processed.push(o);
      if (o === 5 && !crashed) { crashed = true; failed = true; break; }
    }
    if (failed) continue;                       // restart resumes from the committed offset
    if (strategy === 'commit-after') committed = next;
  }
  const seen = processed.reduce((m, o) => ((m[o] = (m[o] ?? 0) + 1), m), {});
  const missing = [...Array(10).keys()].filter((o) => !seen[o]);
  const dupes = Object.keys(seen).filter((o) => seen[o] > 1).map(Number);
  console.log(strategy.padEnd(13), 'missing:', JSON.stringify(missing), 'duplicated:', JSON.stringify(dupes));
}
run('commit-first');   // commit-first  missing: [6,7] duplicated: []
run('commit-after');   // commit-after  missing: [] duplicated: [4,5]
JavaScript

On a cluster, compare the group’s committed offsets with the log end before and after killing a consumer mid-batch:

kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group notifications
# CURRENT-OFFSET is the committed offset; LAG = LOG-END-OFFSET - CURRENT-OFFSET
Terminal

Follow-up questions

Why can lag show zero while work is still pending? Offsets were committed before asynchronous processing finished, as in the worked scenario.

How do you replay a topic for one group? Stop the group, then kafka-consumer-groups.sh --reset-offsets --to-datetime ... --execute, and restart it.

What happens to a group’s offsets if it stops for two weeks? Once the group is empty, its offsets expire after offsets.retention.minutes (7 days) and it falls back to auto.offset.reset on return.

Can you attach information to a commit? Yes, OffsetAndMetadata carries a metadata string, sometimes used to record which application version committed.

Interview exercise

A consumer polls 500 records, processes them in a loop that writes each to PostgreSQL with a plain INSERT, then calls commitSync(). The pod is killed after 300 inserts. What happens on restart, what guarantee does this design give, and what two changes would give exactly-once effects?

Answer and reasoning

Nothing from that batch was committed, so the new owner of the partition resumes at the first of the 500 records and inserts all of them. The first 300 rows are inserted twice, so this is at-least-once with visible duplicates. The first fix is to make the write idempotent: carry an event ID (or use topic, partition and offset) and insert with a unique constraint and ON CONFLICT DO NOTHING. The second, stronger option is to write each row and the partition’s next offset in one database transaction and seek to that stored offset on assignment, so database state and progress can never disagree. Either way, the commit-after-processing order stays: it is what prevents loss.

Continue learning

Practise with the Apache Kafka interview questions and the Kafka MCQs. Related notes: Kafka exactly-once and transactions, Kafka consumer groups and rebalancing and consumer deduplication in microservices. Primary sources: the KafkaConsumer Javadoc for Kafka 4.0, Confluent’s message delivery guarantees and the Apache Kafka documentation.

More in Apache Kafka

esc