A committed offset is the next record a group will read after a restart. With enable.auto.commit (the default) the client commits in the background every auto.commit.interval.ms (5 s). In librdkafka it commits the offsets it has stored, and by default it stores each record's offset as soon as it hands the record to you, before you have processed it. The test reads 1,000 records, processes 200, waits for a background commit and "crashes".
import time
from confluent_kafka import Consumer, TopicPartition
P0 = TopicPartition("booknest.order-events", 0, 0)
for mode in ("auto", "manual"):
group = f"booknest-commit-{mode}-{int(time.time())}"
c = Consumer({"bootstrap.servers": "localhost:33092", "group.id": group,
"auto.commit.interval.ms": 1000, # commit every second
"enable.auto.offset.store": mode == "auto"}) # manual: we store offsets
c.assign([P0])
batch = c.consume(1000, timeout=10)
for msg in batch[:200]: # the job dies after 200 of 1,000 events
if mode == "manual":
c.store_offsets(msg) # mark the event done only after processing it
time.sleep(2) # the background commit runs; then the crash
checker = Consumer({"bootstrap.servers": "localhost:33092", "group.id": group})
done = checker.committed([TopicPartition("booknest.order-events", 0)], timeout=10)
print(f"{mode:<6} processed 200 of {len(batch)}, committed offset {done[0].offset}")auto processed 200 of 1000, committed offset 1000 manual processed 200 of 1000, committed offset 200
With the defaults, 800 unprocessed events would be skipped after a restart: at-most-once. Setting enable.auto.offset.store=false and calling store_offsets() after processing gives at-least-once with the background commit; or disable auto-commit and commit() after each batch. Java's auto-commit only covers records returned by earlier poll() calls.