Skip to content

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

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.

After a crash at 1 s, 1,000 records 4 horizontal bars comparing auto-commit: never processed with the others. After a crash at 1 s, 1,000 records auto-commit: never processed 224 lost auto-commit: processed twice 0 watermark: never processed 0 lost watermark: processed twice 56 redone aiokafka 0.14, one partition, 4 workers at 5 ms per record, process killed after 1.0 s, then restarted. Lost work is silent; duplicate work is handled by idempotency.

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.

Fetch, process, mark, commit A flow of 5 stages. Fetch, process, mark, commit getmany() register offsets workers process concurrently finish(offset) done or dead-lettered watermark contiguous prefix commit every 1 s and on shutdown The commit trails the slowest unfinished record, never the fetch position.

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.

Which commit strategy fits this consumer? A decision on How are records processed with 3 outcomes. Which commit strategy fits this consumer? How are records processed? inline, one at a time auto-commit ok commits trail processing concurrent workers manual low-watermark commits at-least-once effects in your own DB store offsets in the same transaction exactly-once Auto-commit is only safe when the next fetch implies the last record is done.

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.