Kafka Offset Commits with Parallel Workers: Track the Completed Prefix

A tested Python model shows how to commit Kafka offsets safely when workers finish out of order, with offset gaps, bounded queues, and assignment changes.

12 min read

A Kafka consumer hands three records to a worker pool. Record 105 finishes first, so a worker commits offset 106. Then the process crashes. Records 100 and 104 were still running, but the replacement consumer resumes after them.

The commit succeeded. The application lost unfinished work from its normal recovery path.

This constructed failure is easy to introduce when a sequential consumer becomes concurrent. A committed offset represents a resumable position for a partition, not an acknowledgment attached to one independently processed record. The question is therefore not which worker finished most recently. It is which delivered records can safely be left behind on restart.

We will build and test that bookkeeping rule in Python. The example uses no Kafka client or broker: it isolates the completion-order algorithm so every possible ordering of a small workload can be exercised deterministically. Kafka integration requirements are discussed separately, with references to the Apache Kafka 4.1 documentation and its currently served 4.1.2 Java API pages.

A fetched position is not a completed position

Keep three values separate for each topic-partition: the position reached by delivery, the position justified by completed application work, and the position confirmed by an offset commit. They can move at different times.

The KafkaConsumer contract distinguishes consumer position from committed position and documents that offsets can have gaps. A recovery position identifies where reading should resume. It does not tell Kafka whether a database transaction or HTTP operation has completed.

Suppose the application starts at 100, receives offsets 100, 104, and 105, and knows that the position after the batch is 106. While all three are running, the safe recovery position remains 100. Finishing 105 changes nothing. Finishing 100 permits progress to 104. Only when 104 also completes can the already finished 105 join the completed prefix and allow 106.

The gap between 100 and 104 is not evidence that three worker results are missing. Compaction and transactional records can create gaps in the offsets visible to an application. Track the sequence actually delivered; an algorithm waiting for completion notifications for every integer offset can stall permanently.

Define completion before calculating the commit

A worker returning successfully is useful only if success means the required effect has reached its intended durability boundary. Enqueuing an in-memory task, starting a database write, or receiving a future is not necessarily completion.

For this example, a completion notification asserts that the record’s required work has succeeded. The tracker cannot inspect that assertion. If a worker reports success before its database transaction commits, perfectly correct offset bookkeeping still produces an unsafe checkpoint.

Even correctly ordered work and commits leave a crash window: an external write can succeed before the offset commit becomes durable. Recovery can repeat the write. Kafka’s delivery-semantics design discussion explains this distinction and the coordination needed when outputs live outside Kafka.

Kafka transactions can coordinate Kafka output records and consumer offsets within their transaction boundary. They do not automatically make an arbitrary external HTTP call part of that transaction. Decide how repeated external effects are handled before interpreting a successful commit as an end-to-end guarantee.

Track a completed prefix for one assignment

Save the following as progress.py. It was tested with Python 3.13.4 and requires only the standard library.

from collections import deque
from dataclasses import dataclass


@dataclass(frozen=True)
class Ticket:
    assignment: object
    offset: int


class PartitionProgress:
    """Single-owner bookkeeping for one partition assignment; no Kafka I/O."""

    def __init__(self, start_offset: int, capacity: int = 1000):
        if type(start_offset) is not int or start_offset < 0:
            raise ValueError("start_offset must be a nonnegative integer")
        if type(capacity) is not int or capacity <= 0:
            raise ValueError("capacity must be a positive integer")
        self._assignment = object()
        self._active = True
        self._read_position = start_offset
        self._safe_offset = start_offset
        self._capacity = capacity
        self._order = deque()
        self._pending = {}  # offset -> (resume position, completed)

    @property
    def safe_offset(self) -> int:
        if not self._active:
            raise RuntimeError("assignment is no longer active")
        return self._safe_offset

    def offer(self, offsets: list[int], next_position: int) -> list[Ticket]:
        """Register a whole partition batch before dispatching any of its work."""
        if not self._active:
            raise RuntimeError("assignment is no longer active")
        if not offsets or any(type(n) is not int or n < 0 for n in offsets):
            raise ValueError("offsets must be nonempty nonnegative integers")
        if any(a >= b for a, b in zip(offsets, offsets[1:])):
            raise ValueError("offsets must be strictly increasing")
        if offsets[0] < self._read_position:
            raise ValueError("batch overlaps an earlier delivery")
        if type(next_position) is not int or next_position <= offsets[-1]:
            raise ValueError("next_position must follow the batch")
        if len(self._pending) + len(offsets) > self._capacity:
            raise BufferError("pending capacity exceeded; batch not registered")

        resumes = offsets[1:] + [next_position]
        for offset, resume in zip(offsets, resumes):
            self._order.append(offset)
            self._pending[offset] = (resume, False)
        self._read_position = next_position
        return [Ticket(self._assignment, offset) for offset in offsets]

    def complete(self, ticket: Ticket) -> bool:
        """Apply a successful completion on the owner thread."""
        if not self._active or ticket.assignment is not self._assignment:
            return False
        if ticket.offset not in self._pending:
            return False
        resume, done = self._pending[ticket.offset]
        if done:
            return False
        self._pending[ticket.offset] = (resume, True)
        while self._order and self._pending[self._order[0]][1]:
            offset = self._order.popleft()
            self._safe_offset, _ = self._pending.pop(offset)
        return True

    def revoke(self) -> None:
        self._active = False
        self._order.clear()
        self._pending.clear()

There are two contracts behind this small class. First, offer receives every record from a nonempty partition batch, in delivery order, before any of that batch’s work is dispatched. Its next_position comes from the consumer’s corresponding batch metadata, not the highest offset observed somewhere else.

Second, one owner thread calls every method. Workers return tickets through a thread-safe completion queue; they do not mutate the tracker directly. The object is deliberately a state machine, not a concurrent collection or Kafka adapter.

Each pending entry stores the position to resume after it. For interior records, that position is the next delivered offset. For the last record, it is the supplied position after the batch. Completed entries are removed only from the front. A completed tail remains pending until every earlier delivered record succeeds.

Registration validates the whole batch before changing state. A capacity failure therefore cannot leave half a batch registered. The caller still owns that rejected batch and must retain it in bounded storage or recover it through an explicit, correctly coordinated replay path. Dropping it and continuing with later batches violates the input contract.

With valid input, the offset never advances past an unfinished registered record. It can remain conservatively behind a later visible record when there is a gap between batches. That is safe: a checkpoint need not claim the largest conceivable position to be useful.

Reproduce the out-of-order completion

Save this as demo.py in the same directory and run python3 demo.py.

from progress import PartitionProgress

progress = PartitionProgress(start_offset=100, capacity=3)
a, b, c = progress.offer([100, 104, 105], next_position=106)
print(f"initial: safe={progress.safe_offset}")
for ticket in (c, a, b):
    progress.complete(ticket)
    print(f"completed {ticket.offset}: safe={progress.safe_offset}")

old_ticket = a
progress.revoke()
replacement = PartitionProgress(start_offset=100)
print(f"old completion accepted: {replacement.complete(old_ticket)}")

The program produces:

initial: safe=100
completed 105: safe=100
completed 100: safe=104
completed 104: safe=106
old completion accepted: False

The first completion is the tempting mistake. Committing 106 at that point would abandon two unfinished records on recovery. The tracker holds at 100 instead. After 100 completes, it advances directly to the next delivered record, 104, without inventing work for the missing integer offsets.

The final line exercises a different boundary. A completion ticket belongs to the assignment that created it. A replacement tracker rejects it even if the same offset is being processed again. The opaque identity distinguishes local work lifetimes; it is not Kafka’s group generation, a broker leader epoch, or a distributed fencing token.

If the application seeks backward or otherwise changes its delivery history, retire the tracker and create a new one for the new execution context. Reusing its pending entries across a seek would mix two different attempts at processing the same offsets.

Test the invariant across every completion order

Timing-based tests can accidentally exercise the same worker order on every run. For five records, there are only 120 permutations, so test them all. Save this as test_article.py and run python3 -m unittest test_article -v.

import itertools
import unittest
from progress import PartitionProgress


class PrefixInvariant(unittest.TestCase):
    def test_every_completion_order(self):
        offsets = [100, 104, 105, 109, 115]
        for order in itertools.permutations(range(len(offsets))):
            p = PartitionProgress(100)
            tickets = p.offer(offsets, next_position=120)
            completed = set()
            for index in order:
                p.complete(tickets[index])
                completed.add(index)
                unfinished = [i for i in range(len(offsets))
                              if i not in completed]
                expected = offsets[unfinished[0]] if unfinished else 120
                self.assertEqual(p.safe_offset, expected)


if __name__ == '__main__':
    unittest.main()

The expected position is calculated from the first unfinished item in the original delivery sequence, independently of the tracker’s queue operations. The test checks after every completion, producing 600 intermediate assertions. It includes nonconsecutive offsets and a final resume position beyond the last record.

The wider validation suite also checks capacity rejection without partial mutation, completed entries blocked behind a slow head, multiple batches sharing one frontier, duplicate notifications, invalid input, independent partitions, and old assignments. All 11 tests passed. These are tests of the bookkeeping model; they do not exercise broker failover, a client library’s callbacks, or external write durability.

Keep commit ownership in the consumer loop

For a Java KafkaConsumer integration, keep polling and client operations on the consumer’s owning thread; that client is not thread-safe. Dispatch work outward and bring completion notifications back. Disable automatic commits when the application independently determines completed progress. These constraints follow the manual-control and threading guidance in the consumer API documentation.

The owner maintains one tracker per active topic-partition assignment. Before producing a commit map, it drains a bounded amount of completion work, verifies ownership, and reads each partition’s safe position. Never use the maximum offset across partitions: offset 500 in one partition says nothing about offset 20 in another.

The integer in the example is not a complete Java commit object. Preserve the applicable leader-epoch metadata when adapting it to OffsetAndMetadata. The Java client documents using the next record’s offset, or the batch’s ConsumerRecords.nextOffsets() when exhausted. Carry the corresponding metadata through the adapter rather than reconstructing it from unrelated current state.

Keep a separate record of confirmed commits. Advancing safe_offset means the application may propose that checkpoint; it does not mean the broker has accepted it. A failed or uncertain commit must not update a dashboard field labeled “confirmed.”

A simple integration can serialize explicit commits on the owner thread. An asynchronous design needs an explicit policy for requests in flight, callback results, and assignment changes. Do not blindly replay an old commit map from a callback after ownership has changed. Re-evaluate current ownership and current safe progress before any application-level retry.

A slow head is a capacity problem too

Imagine the first record waits on a failing dependency while thousands of later records complete. Their work may be done, but the recovery position cannot pass the first record. If the process crashes, those later effects may all be repeated.

The capacity bound counts all retained entries, including completed ones behind the gap. Counting only actively running workers would miss the growing completion backlog. Bound worker queues and payload bytes separately: this tracker bounds record metadata, not all memory held by the application or client.

Pause delivery from a pressured partition before its available capacity is exhausted, accounting for an entire batch that may be returned. Continue polling as required for group participation and service other partitions when possible. Ensure the completion queue cannot block the owner in a circular wait with workers trying to report results.

The consumer configuration reference specifies that max.poll.records limits records returned by a poll, not the underlying fetch cache. It also documents max.poll.interval.ms and the separate behavior of static members. Neither setting replaces an application admission policy or a byte budget.

A safe commit frontier also does not restore business ordering. If two workers update the same account out of order, withholding the commit cannot undo that ordering error. Use sequential partition processing, an appropriate per-key execution policy, or business operations whose correctness permits concurrency.

Revocation must invalidate old work locally

During orderly revocation, stop new dispatch for the affected partitions. If the integration allows a bounded drain, incorporate only successful completions observed within that ownership window. Commit only the resulting safe positions while the client still permits it, then retire the trackers. Do not wait indefinitely for a stuck worker.

The rebalance listener contract distinguishes revocation from partition loss. With onPartitionsLost, another member may already own the partitions. Handle loss as loss of authority, not as an ordinary opportunity to commit. Cooperative changes can affect only a subset, so keep valid state for assignments that remain.

Calling revoke() invalidates this model’s tickets and checkpoint accessor. It does not cancel a running thread, stop an external write, or revoke a checkpoint integer that somebody already copied. The owner must prevent stale commit requests from escaping, and the destination may still need idempotency or actual fencing to handle late effects.

After reassignment, initialize progress from the chosen recovery position and issue fresh tickets. Do not import the old tracker’s completed tail merely because its offsets match: its work and commit history belong to an earlier ownership interval.

Observe the gap that prevents recovery progress

Track the age and offset of the earliest unfinished record, retained entry count, completed-but-blocked count, safe position, confirmed position, and time spent paused. Those signals distinguish slow processing from a commit failure and reveal when one record is holding a large replay window open.

Offset distance is useful as a position gap, but it is not an exact record count when offsets have holes. Use the tracker’s entry counts for retained application work. For latency, measure elapsed time from delivery to completion and from completion to confirmed checkpoint.

A permanently failing record needs a stated policy: retry with bounds, stop for intervention, or complete a durable error-handling workflow before considering it resolved. Marking it complete solely to make the lag graph fall turns an operational problem into an unrecorded loss.

The commit frontier answers one narrow question: which delivered work can this assignment safely leave behind? Keeping that answer explicit makes concurrency reviewable. Recovery semantics, business ordering, and external effects still need their own evidence.

What do you think?

Add your perspective.

Your email address will not be published. Required fields are marked *

What is on your mind?

START WITH A TOPIC
SAVED FOR A QUIET MOMENT

My reading list

Your list is stored only in this browser.

See you in the next story.

New ideas, new stories. The same curiosity.

Open the RSS feed