Skip to content

Stopping Kafka Consumers Gracefully

When a Kafka consumer stops, its partitions have to move to another member of the group, and how it stops decides how long they sit idle and how much work is done twice. Measured on Python 3.14 with aiokafka 0.14.0 against Kafka 4.1.0, two consumers in one group sharing a 4-partition topic, each message taking 20 ms to process and offsets committed after every batch: killing one consumer with SIGKILL left its partitions unprocessed until the broker's 10-second session timeout expired — the other consumer received all four partitions 12.1 s after the kill — and 89 messages were processed twice. Stopping it gracefully — finishing the batch, committing, calling consumer.stop() to leave the group — handed the partitions over in 3.1 s. But the graceful stop still produced 50–100 duplicates, from the other consumer: the rebalance revoked its partitions in the middle of its own batch, and its commit after the batch did not count. Committing processed offsets in an on_partitions_revoked listener and stopping a batch when its partition is revoked brought that down to 6. This guide builds the consumer that stops cleanly and keeps its neighbours clean too.

Prerequisites

1. Measure what a crash costs

A consumer that commits after each batch, stopped without any shutdown handling:

consumer = AIOKafkaConsumer(bootstrap_servers=BS, group_id="orders",
                            enable_auto_commit=False, max_poll_records=50)
consumer.subscribe(["orders"])
await consumer.start()
while True:
    batch = await consumer.getmany(timeout_ms=200, max_records=50)
    for tp, messages in batch.items():
        for m in messages:
            await handle(m)                       # 20 ms each
    if batch:
        await consumer.commit()

Measured by killing it with SIGKILL while a second consumer ran: the broker kept the dead member in the group until session_timeout_ms — 10,000 ms by default in aiokafka — expired, and the survivor received all four partitions 12.1 s after the kill. For those 12 seconds, half the topic was not consumed at all. Then the survivor re-read from the last committed offsets, processing 89 messages a second time — the dead consumer's uncommitted batch. A plain SIGTERM has the same effect if nothing handles it: the default action ends the process just as abruptly.

Verify: a test that kills a consumer measures how long its partitions are idle and how many messages are processed twice.

Two consumers, 4 partitions, one stopped mid-stream A grid of 3 rows by 3 columns. Two consumers, 4 partitions, one stopped mid-stream stop method partitions reassigned after processed twice SIGKILL (crash) 12.1 s 89 SIGTERM: finish batch, commit, consumer.stop() 3.1 s 50-100 graceful + commit on revoke, stop batch when revoked 3.1 s 6 aiokafka 0.14.0, Kafka 4.1.0; 20 ms per message, batches of 50.

2. Stop fetching, finish the batch, leave the group

Handle SIGTERM by setting an event the loop checks between batches, then commit and stop the consumer so it sends a LeaveGroup request:

async def run(consumer):
    stop = asyncio.Event()
    asyncio.get_running_loop().add_signal_handler(signal.SIGTERM, stop.set)
    await consumer.start()
    try:
        while not stop.is_set():
            batch = await consumer.getmany(timeout_ms=200, max_records=50)
            for tp, messages in batch.items():
                for m in messages:
                    await handle(m)
            if batch:
                await consumer.commit()
    finally:
        await consumer.stop()                     # leaves the group: no session-timeout wait

Measured: the stopping consumer exited 0.18–0.25 s after SIGTERM, and its partitions were assigned to the other consumer 3.1 s later — against 12.1 s after a kill. Bound the batch so the drain fits the platform's grace period: max_records times the per-message time is the worst case, here 50 × 20 ms = 1 s.

Verify: after SIGTERM, the consumer exits within one batch's processing time, and the group rebalances without waiting for the session timeout.

3. Protect the other consumers' batches

A graceful stop still triggers a rebalance, and with aiokafka's default eager protocol, every member's partitions are revoked and reassigned — including those of consumers that are not stopping. Measured: the graceful version caused 50–100 duplicates after the stop, all in the surviving consumer. It had been in the middle of a batch when its partitions were revoked; its commit after the batch was not applied to the new assignment, and it re-read the batch. Commit what has actually been processed when partitions are revoked, and stop processing a partition you no longer own:

processed: dict[TopicPartition, int] = {}

class CommitOnRevoke(ConsumerRebalanceListener):
    def __init__(self, consumer):
        self.consumer = consumer

    async def on_partitions_revoked(self, revoked):
        offsets = {tp: OffsetAndMetadata(processed[tp], "") for tp in revoked if tp in processed}
        if offsets:
            await self.consumer.commit(offsets)
        for tp in revoked:
            processed.pop(tp, None)

    async def on_partitions_assigned(self, assigned):
        pass

# in the processing loop
for tp, messages in batch.items():
    for m in messages:
        if tp not in consumer.assignment():
            break                                 # revoked mid-batch: leave the rest
        await handle(m)
        processed[tp] = m.offset + 1

Measured: duplicates after the stop fell to 6, with the same 3.1-second handover. A handful remain because a message can finish processing in the moment between revocation and its offset being recorded; handlers still need to be idempotent.

Verify: a test that stops one consumer during the other's batch shows only a handful of duplicates, and none of them cause incorrect state.

Graceful stop with commit on revoke A sequence of 6 messages between 3 participants. Graceful stop with commit on revoke consumer A coordinator consumer B SIGTERM: finish batch, commit LeaveGroup (consumer.stop()) rebalance: revoke B's partitions on_partitions_revoked: commit processed offsets assign all 4 partitions (3.1 s) resume from committed offsets: 6 duplicates The stopping consumer and the surviving one both commit before letting go.

4. Bound the stop within the platform's deadline

Graceful shutdown must still end. The drain is at most one batch, but consumer.stop() talks to the broker and can wait on an unreachable one. Put a deadline around the whole sequence:

async def shutdown(consumer, grace: float = 10.0):
    try:
        async with asyncio.timeout(grace):
            await consumer.commit()
            await consumer.stop()
    except (TimeoutError, KafkaError) as e:
        log.warning("consumer did not stop cleanly: %s", e)   # the session timeout will clean up

If the deadline passes, the worst case is the crash case from step 1 — partitions idle until the session timeout, one batch reprocessed — which idempotent handlers already tolerate. Keep the grace below the platform's termination period, and add the process-level backstop from enforcing a hard shutdown deadline.

Verify: with the broker unreachable, SIGTERM still ends the process within the grace period.

5. Tune the session timeout for crashes

Graceful stops avoid the session timeout; crashes depend on it. A shorter session_timeout_ms makes the group notice a dead consumer sooner, at the cost of false evictions when a consumer is briefly slow to heartbeat:

consumer = AIOKafkaConsumer(
    bootstrap_servers=BS, group_id="orders",
    session_timeout_ms=10_000,                # crash detection: idle partitions up to ~10 s
    heartbeat_interval_ms=3_000,              # about a third of the session timeout
    max_poll_interval_ms=300_000,             # longest allowed gap between getmany calls
)

Measured with the default 10 s, the crashed consumer's partitions sat idle for 12.1 s. aiokafka heartbeats from a background task, so a slow handler does not miss heartbeats — but a handler that blocks the event loop does, and then a too-short session timeout evicts healthy consumers. For partition planning and per-key concurrency inside each consumer, see scaling Kafka consumers with partitions.

Verify: the session timeout is chosen deliberately, and the consumer's loop never blocks for longer than a heartbeat interval.

Seconds until the survivor owned all partitions 2 horizontal bars comparing SIGKILL: wait for the 10 s session timeout with the others. Seconds until the survivor owned all partitions SIGKILL: wait for the 10 s session timeout 12.1 s graceful: LeaveGroup via consumer.stop() 3.1 s The difference is the session timeout, which a LeaveGroup skips.

Verification

Consumers stop gracefully when:

  • SIGTERM ends fetching, the current batch finishes, offsets are committed, and consumer.stop() leaves the group.
  • Rebalance listeners commit processed offsets on revocation, and batches stop on revoked partitions.
  • The whole stop is bounded below the platform's grace period.
  • Handlers are idempotent, since a few duplicates remain even in the best case.

Diagnostic Hook: when consumer lag spikes on half the partitions for about ten seconds after every deploy, the stopping consumers are not leaving the group. A killed consumer's partitions were reassigned after 12.1 s here, against 3.1 s after consumer.stop().

Pitfalls & edge cases

  • No SIGTERM handling. Partitions idle for the session timeout; measured 12.1 s.
  • Committing only at the end of each batch. Surviving consumers re-read 50-100 messages.
  • Continuing a batch after revocation. Another consumer processes the same messages.
  • An unbounded consumer.stop(). An unreachable broker can stall shutdown.

Frequently Asked Questions

How do I stop an aiokafka consumer gracefully?

On SIGTERM, stop fetching, finish the current batch, commit, and await consumer.stop(), which leaves the group. Partitions moved in 3.1 s instead of 12.1 s.

Why does Kafka take 10 seconds to reassign partitions after a consumer dies?

The broker waits for session_timeout_ms (10 s by default in aiokafka) before evicting a member that did not leave. A graceful stop skips that wait.

Why do other consumers reprocess messages when one stops?

The rebalance revokes their partitions mid-batch and their end-of-batch commit is lost. Committing in on_partitions_revoked cut duplicates from 50-100 to 6.

Can a graceful Kafka shutdown avoid all duplicates?

Not entirely: 6 duplicates remained with commit-on-revoke. Keep handlers idempotent.