Command Palette

Search for a command to run...

Hectal
PHASE 3Intermediate ~9 min· topic 2 of 5

Topic 3.2

Offsets, Commits, and Lag

In one line

A committed offset is the position of the next record the group should read, stored in __consumer_offsets. Commit after processing for at-least-once; commit before for at-most-once. Committing an offset means "everything before this is done", so you can't skip a failed record and commit past it without deciding what happens to it.

0/5 · 0%

Think of it like this

A bookmark. Placing it at page 104 means "I've finished everything up to 103". If you skipped page 101 because it was unreadable, a bookmark at 104 means you'll never come back to 101 unless you wrote it down somewhere else.

Key ideas

  1. 01

    Positions: current position (next offset poll returns), committed offset (next offset to read after a restart, stored per group and partition), log end offset (LEO, the next offset to be written), and high watermark (last offset replicated to all ISR, the limit consumers can read). Lag = LEO (or high watermark) − committed offset.

  2. 02

    Auto commit (enable.auto.commit=true, every 5 s): commits the positions returned by the previous poll during the next poll. Records are usually processed before being committed, but a crash between process and auto-commit reprocesses up to 5 s of records, and processing records asynchronously after poll can commit before they're done.

  3. 03

    Manual commit: commitSync() blocks and retries until success or a fatal error; commitAsync() doesn't block and doesn't retry (a later commit supersedes it). Common pattern: commitAsync in the loop, commitSync on shutdown and on partition revocation.

  4. 04

    Commit specific offsets with commitSync(Map<TopicPartition, OffsetAndMetadata>) and commit offset + 1: the committed value is the next record to read.

  5. 05

    Offset retention: committed offsets for a group are kept for offsets.retention.minutes (7 days) after the group becomes empty; after that a restarted group falls back to auto.offset.reset.

Code & diagrams

offsets.txttext
partition 3:  0  1  2  3  4  5  6  7  8  9  10 11
                            ^           ^           ^
                   committed=4    position=7   log end offset=12
   records 4..6 were returned by poll and are being processed
   lag = 12 - 4 = 8
per-record-commit.javajava
Map<TopicPartition, OffsetAndMetadata> toCommit = new HashMap<>();
for (ConsumerRecord<String, String> r : records) {
    handle(r);                                                    // may throw
    toCommit.put(new TopicPartition(r.topic(), r.partition()),
                 new OffsetAndMetadata(r.offset() + 1));          // NEXT offset to read
}
consumer.commitAsync(toCommit, (offsets, ex) -> {
    if (ex != null) log.warn("async commit failed, a later commit will cover it", ex);
});

Interview problem

The problem

Offsets 100, 101, 102: 101 fails

A poll returns offsets 100, 101 and 102. Processing gives 100 success, 101 failure, 102 success. When should the consumer commit? What happens if it commits offset 103?

The interviewer follows up

01

What if the consumer crashes after processing 100–102 but before committing?

When it breaks

Auto commit with asynchronous processing after poll

What you see

Offsets are committed for records still being processed in background threads; a crash loses them.

Fix & prevent

Disable auto commit and commit only offsets whose records are fully processed.

Explain it without notes

01

Define current position, committed offset, log end offset and lag.

Practice

01

Produce 10 records, consume 5 with manual commits, stop, and check lag with kafka-consumer-groups --describe. Then reset the group to the earliest offset.

Trade-offs

  • ↔

    Frequent commits reduce reprocessing after crashes but add broker load; batch commits are cheaper but reprocess more.

Done when you can

  • I commit offset + 1 after processing and never commit past an unhandled failure.