Committing Kafka Offsets Safely in Async Consumers¶
A Kafka consumer's committed offset is its promise that everything before it has been handled. With auto-commit, the client commits periodically whatever it has fetched, not what it has processed — harmless when processing happens inline before the next fetch, and a data-loss bug when an async consumer hands records to concurrent workers. Tested with aiokafka 0.14: a consumer that fed 1,000 records to four workers through a queue, with auto-commit every 100 ms, crashed after one second; on restart the group resumed from the auto-committed offset and 224 messages were never processed. The same consumer committing a low watermark — the highest offset below which every record had been processed — lost none after the same crash and reprocessed 56, the records that were in flight. This guide implements watermark commits, handles rebalances, and makes the redo harmless.
Prerequisites¶
- Python 3.11+,
pip install aiokafka; measured with aiokafka 0.14. - Consumer basics, from consuming Kafka topics with aiokafka.
- Idempotency, from making background jobs idempotent.
1. See how auto-commit loses messages¶
Auto-commit commits the position of the last fetched record on a timer. When fetching and processing are decoupled, the commit runs ahead of the work:
consumer = AIOKafkaConsumer("orders", group_id="billing", enable_auto_commit=True,
auto_commit_interval_ms=100)
await consumer.start()
work: asyncio.Queue = asyncio.Queue()
async for record in consumer:
work.put_nowait(record) # fetched -> committed soon, processed later
# workers take records from `work` and process them concurrently
Measured: 1,000 records, four workers taking 5 ms each, crash after 1 s. By then the consumer had fetched — and auto-committed — far more than the workers had finished. The restarted consumer began at the committed offset, and 224 records that had been fetched but not processed were skipped forever. No error was raised anywhere; the loss is only visible by counting.
Verify: compare the committed offset with the offset of the last record your workers finished; with auto-commit and a work queue, the committed offset is ahead.
2. Disable auto-commit and track a low watermark¶
Commit only offsets whose records — and every record before them — are finished. With concurrent workers, records finish out of order, so keep the set of finished offsets per partition and advance the watermark over the contiguous prefix:
from aiokafka import AIOKafkaConsumer, TopicPartition
class Watermarks:
def __init__(self) -> None:
self.next: dict[TopicPartition, int] = {} # first offset not yet known done
self.done: dict[TopicPartition, set[int]] = {}
def start(self, tp: TopicPartition, offset: int) -> None:
self.next.setdefault(tp, offset)
self.done.setdefault(tp, set())
def finish(self, tp: TopicPartition, offset: int) -> bool:
self.done[tp].add(offset)
moved = False
while self.next[tp] in self.done[tp]:
self.done[tp].remove(self.next[tp])
self.next[tp] += 1
moved = True
return moved # True when there is something new to commit
consumer = AIOKafkaConsumer("orders", group_id="billing", enable_auto_commit=False,
auto_offset_reset="earliest")
Kafka's committed offset is the offset of the next record to read, which is why the watermark stores "first offset not yet done". A slow record holds the watermark back, and everything after it is redone on a crash — that is the measured 56 duplicates, and it is the price of losing nothing. The same idea for job checkpoints is in checkpointing progress in long-running async jobs.
Verify: in a test with random processing delays, the committed offset never exceeds the oldest unfinished record's offset.
3. Commit periodically, off the hot path¶
Committing after every record is a broker round trip per record. Commit the watermarks on a short timer, and once more on shutdown:
async def committer(consumer, marks: Watermarks, every: float = 1.0) -> None:
last: dict[TopicPartition, int] = {}
while True:
await asyncio.sleep(every)
offsets = {tp: n for tp, n in marks.next.items() if last.get(tp) != n}
if offsets:
await consumer.commit(offsets)
last.update(offsets)
async def worker(consumer, work: asyncio.Queue, marks: Watermarks) -> None:
while True:
record = await work.get()
tp = TopicPartition(record.topic, record.partition)
try:
await process(record)
finally:
marks.finish(tp, record.offset) # done or dead-lettered: either way, finished
work.task_done()
The interval bounds how much is redone after a crash: at most what was processed since the last commit, plus whatever was in flight. A record that failed permanently must still be finished — after sending it to a dead-letter topic — or it holds the watermark back forever; see handling poison messages in async consumers.
Verify: the consumer group's committed offsets advance every interval under load, and lag measured from them matches the work queue depth plus in-flight records.
4. Bound the work in flight¶
Without a bound, the fetch loop keeps pulling records into memory while workers lag, which grows both memory and the amount redone after a crash. Bound the queue and pause fetching when it is full:
async def run(consumer, marks: Watermarks, workers: int = 8, max_in_flight: int = 500) -> None:
work: asyncio.Queue = asyncio.Queue(maxsize=max_in_flight)
async with asyncio.TaskGroup() as tg:
for _ in range(workers):
tg.create_task(worker(consumer, work, marks))
tg.create_task(committer(consumer, marks))
while True:
batches = await consumer.getmany(timeout_ms=500, max_records=100)
for tp, records in batches.items():
for record in records:
marks.start(tp, record.offset)
await work.put(record) # blocks when max_in_flight are queued
await work.put(...) stops the loop from fetching more until workers catch up. Keep max_in_flight modest: the loop must still call getmany often enough to stay within max.poll.interval.ms, or the broker considers the consumer dead and rebalances. Per-key ordering, if you need it, comes from routing records with the same key to the same worker, as in processing queue items in order per key.
Verify: consumer memory and in-flight count stay flat when processing is artificially slowed.
5. Commit on rebalance and make processing idempotent¶
When partitions move to another consumer, commit what is finished for them first, or the new owner will redo more than necessary — or, with careless code, the old owner will commit an offset for a partition it no longer owns:
from aiokafka.abc import ConsumerRebalanceListener
class CommitOnRevoke(ConsumerRebalanceListener):
def __init__(self, consumer, marks: Watermarks) -> None:
self.consumer, self.marks = consumer, marks
async def on_partitions_revoked(self, revoked) -> None:
offsets = {tp: self.marks.next[tp] for tp in revoked if tp in self.marks.next}
if offsets:
await self.consumer.commit(offsets)
for tp in revoked:
self.marks.next.pop(tp, None)
self.marks.done.pop(tp, None)
async def on_partitions_assigned(self, assigned) -> None:
pass
consumer.subscribe(["orders"], listener=CommitOnRevoke(consumer, marks))
Records still in flight for a revoked partition will be redone by the new owner. That redo — like the 56 duplicates after a crash — is unavoidable with at-least-once delivery, so processing must be idempotent: upserts keyed by a business id, a processed-message table checked in the same transaction, or naturally idempotent operations. With idempotent handlers, the watermark design loses nothing and its duplicates are harmless.
Verify: trigger a rebalance by starting a second consumer mid-run; the total of unique processed records is complete and duplicates are absorbed by the handler.
Verification¶
Offset commits are safe when:
- Auto-commit is off for any consumer that processes records concurrently.
- The committed offset is a low watermark of finished records per partition.
- In-flight work is bounded, and commits happen periodically and on revoke.
- Processing is idempotent, so redone records are harmless.
Diagnostic Hook: periodically compare, per partition, the committed offset with the oldest unfinished offset in the process. A committed offset ahead of unfinished work is the auto-commit bug; a watermark stuck for a long time while others advance is one stuck record — log its offset and age.
Pitfalls & edge cases¶
- Auto-commit with a work queue. Measured: 224 of 1,000 records lost in one crash.
- Committing the highest finished offset. Earlier unfinished records are skipped.
- Never finishing failed records. The watermark stops and lag grows forever.
- Non-idempotent handlers. Redone records cause double effects.
Frequently Asked Questions¶
Is Kafka auto-commit safe with asyncio consumers?
Not when records are processed concurrently after being fetched. Auto-commit commits fetched positions, so a crash skips fetched-but-unprocessed records: 224 of 1,000 in testing.
How do I commit Kafka offsets with concurrent processing?
Disable auto-commit, record each finished offset per partition, and periodically commit the low watermark: the first offset that is not yet finished.
Why do I get duplicate messages after a consumer restart?
Records processed after the last commit, or still in flight at the crash, are redelivered. In testing, a watermark-committing consumer redid 56 records and lost none. Make processing idempotent.
What should a consumer do when partitions are revoked?
Commit the watermarks for the revoked partitions in on_partitions_revoked and drop their tracking state; in-flight records for them will be redone by the new owner.
Related¶
- Message Brokers & Event Streams — up to the topic overview.
- Publishing to Kafka with aiokafka producers — the producing side of the topic.
- Network I/O & Protocol Handling — the section overview.