Skip to content

Repository files navigation

faststream-concurrent-aiokafka

PyPI version Supported Python versions Downloads Coverage CI License GitHub stars Context7 uv Ruff ty

Concurrent message processing middleware for FastStream with aiokafka.

By default FastStream processes Kafka messages sequentially, one message at a time per subscriber. This library turns each incoming message into an asyncio task so multiple messages are handled concurrently, while keeping offset commits correct and shutdown graceful.

Features

  • Concurrent message processing via asyncio tasks
  • Configurable concurrency limit (semaphore-based)
  • Batch offset committing per partition after each task completes
  • Rebalance-safe: pending offsets are flushed on partition revocation via ConsumerRebalanceListener
  • Fast shutdown: cancels in-flight tasks; uncommitted offsets are redelivered on restart (at-least-once)
  • Does not register SIGTERM/SIGINT handlers; stop it from your lifespan.
  • Handler exceptions are logged but do not crash the consumer
  • Health check helper to probe handler status from a ContextRepo

Installation

pip install faststream-concurrent-aiokafka

Quick start

ack_policy=AckPolicy.MANUAL is required for concurrent processing. Without it, FastStream would commit offsets before processing tasks complete, causing silent message loss on crash. Subscribers on any other ack policy (ACK_FIRST, ACK, REJECT_ON_ERROR, NACK_ON_ERROR) pass through without concurrent processing and behave exactly as they would if this middleware were not registered, so a single broker-level registration is safe across a mix of subscribers.

With AsgiFastStream, the lifespan receives an app-level ContextRepo separate from broker.context. Pass broker.context explicitly instead of the injected argument.

from contextlib import asynccontextmanager
from faststream import ContextRepo
from faststream.asgi import AsgiFastStream
from faststream.kafka import KafkaBroker
from faststream.middlewares import AckPolicy
from faststream_concurrent_aiokafka import (
    KafkaConcurrentProcessingMiddleware,
    initialize_concurrent_processing,
    stop_concurrent_processing,
)

broker = KafkaBroker(...)
# Register KCM on the broker before any other middleware (see DI note below)
broker.add_middleware(KafkaConcurrentProcessingMiddleware)


@asynccontextmanager
async def lifespan(_context: ContextRepo):
    await initialize_concurrent_processing(
        context=broker.context,
        concurrency_limit=20,  # max concurrent tasks (minimum: 1)
        commit_batch_size=100,  # commit after this many pending tasks
        commit_batch_timeout_sec=5.0,  # or after this many seconds
    )
    try:
        yield
    finally:
        await stop_concurrent_processing(broker.context)


app = AsgiFastStream(broker, lifespan=lifespan)


@broker.subscriber("my-topic", group_id="my-group", ack_policy=AckPolicy.MANUAL)
async def handle(msg: str) -> None: ...


# Any non-MANUAL policy (ACK_FIRST is FastStream's default) is passed through
# unchanged, not processed concurrently
@broker.subscriber("other-topic", group_id="other-group", ack_policy=AckPolicy.ACK_FIRST)
async def handle_other(msg: str) -> None: ...

Core concepts

KafkaConcurrentHandler and KafkaBatchCommitter are internal: they are not exported from the package and are described here only to explain the behavior.

KafkaConcurrentProcessingMiddleware

A FastStream BaseMiddleware subclass. Add it to your broker to enable concurrent processing. It wraps each incoming message in an asyncio task submitted to KafkaConcurrentHandler.

KafkaConcurrentHandler

The processing engine. It manages:

  • An asyncio.Semaphore to enforce concurrency_limit
  • In-flight task tracking via a set[asyncio.Task]; each task's done-callback releases the semaphore, removes the task from the set, and logs any non-cancellation exception at ERROR with a traceback
  • FastStream control signals raised by a middleware registered after this one are absorbed before they can end the task, so they neither pin the message body via a traceback nor reach error reporters that wrap asyncio tasks. See Limitations for which are honoured and which only log
  • A KafkaBatchCommitter for offset commits
  • A ConsumerRebalanceListener on every concurrent subscriber that flushes pending commits when partitions are revoked. initialize_concurrent_processing attaches it automatically, including to subscribers declared on routers; see Rebalance handling

This library does not install signal handlers. Shutdown is driven by your lifespan or process manager calling stop_concurrent_processing.

KafkaBatchCommitter

Runs as a background asyncio task. A streaming loop absorbs KafkaCommitTask objects into per-partition pending state and commits each partition's contiguous-done prefix when total pending crosses commit_batch_size, when commit_batch_timeout_sec fires, or when commit_all/close sets the flush event. Cancelled tasks are a hard boundary: the offset advance stops at the cancelled task so it gets redelivered on restart (at-least-once). If the committer's task dies, CommitterIsDeadError is raised to callers.

API reference

initialize_concurrent_processing(context, ...)

Create and start the concurrent processing handler; store it in FastStream's context; attach a rebalance listener to every concurrent subscriber. Call it before the broker starts.

Parameter Default Description
context required FastStream ContextRepo instance
concurrency_limit 10 Max concurrent asyncio tasks (minimum: 1)
commit_batch_size 10 Number of pending tasks that triggers a commit; a commit can include more
commit_batch_timeout_sec 10.0 Max seconds before flushing a batch
shutdown_timeout_sec 20.0 Max seconds the batch committer waits for its background task to drain before forcing cancellation
max_uncommitted_tasks 10000 Max tasks accepted but not yet committed before the consume path blocks (backpressure). None disables the bound.
rebalance_flush_timeout_sec 10.0 Max seconds the attached rebalance listener waits for in-flight tasks when partitions are revoked

Returns the KafkaConcurrentHandler instance.

Each uncommitted entry holds only commit metadata (a task reference, its TopicPartition, offset, and consumer reference), not the message payload, so the default max_uncommitted_tasks of 10000 is on the order of a few MB. Lower it to tighten the memory bound during a commit or broker outage, at the cost of stalling consumption sooner. Keep it >= commit_batch_size so size-based batching can still trigger (below that, commits fall back to the timeout/flush path); set it to None to disable the bound and allow unbounded buffering.

stop_concurrent_processing(context)

Cancel all in-flight tasks, flush completed offsets via the committer, then stop the handler. Uncommitted offsets (from cancelled tasks or anything queued past a cancelled offset) are redelivered on restart (at-least-once).

is_kafka_handler_healthy(context)

Returns True if the KafkaConcurrentHandler stored in context is running and healthy, False otherwise (not initialized, stopped, or committer task dead). Useful for readiness/liveness probes.

KafkaConcurrentProcessingMiddleware

FastStream middleware class. Register it via broker.add_middleware(...). See Quick start for usage examples.

This middleware must be the outermost one. consume_scope fires the handler as a background task and returns None immediately, so any middleware that wraps it on the outside sees that premature return and misfires: wrong timing, early cleanup, or missed exceptions. Middlewares added after it (inner in the chain) run correctly inside the background task.

DI framework compatibility (modern-di-faststream and similar)

DI frameworks like modern-di-faststream register a broker-level middleware that creates a REQUEST-scoped dependency container around each message. If that middleware is outer to KafkaConcurrentProcessingMiddleware, its scope closes as soon as consume_scope returns (before the background task runs), so any dependencies resolved inside the task (database sessions, repositories, …) are created from an already-closed container. Their finalizers never run, leaving connections unreturned to the pool.

To avoid this, call broker.add_middleware(KafkaConcurrentProcessingMiddleware) before setup_di(...) (or any equivalent DI bootstrap call). FastStream stacks broker middlewares so the first registered is outermost. Adding KCM first makes it wrap the DI middleware, so the DI middleware runs inside KCM's background task and can manage the scope lifetime correctly.

broker = KafkaBroker(...)
broker.add_middleware(KafkaConcurrentProcessingMiddleware)  # registered first → outermost
modern_di_faststream.setup_di(app, container=container)  # registered after → inner to KCM

How it works

Message dispatch

On each incoming message, consume_scope calls handle_task(), which acquires a semaphore slot then fires the handler coroutine as a background asyncio.Task.

Concurrency control

The semaphore blocks new tasks when concurrency_limit is reached. The slot is released via a done-callback when the task finishes or fails.

Offset committing

Each dispatched task is paired with its Kafka offset and consumer reference and enqueued in KafkaBatchCommitter. Once the task completes, the committer groups offsets by partition and calls consumer.commit(partitions_to_offsets) with offset + 1 (Kafka's "next offset to fetch" convention).

Rebalance handling

initialize_concurrent_processing attaches a ConsumerRebalanceListener to every subscriber it processes concurrently: AckPolicy.MANUAL, not batch=True, subscribed by topic or pattern, on any Kafka broker of the FastStream application in the context, including subscribers from included routers. It must run before the broker starts (in the lifespan, as above), because FastStream hands the listener to aiokafka when the consumer subscribes; if the broker has already started, or the context holds no FastStream application, it logs an ERROR and attaches nothing. A listener= you pass yourself is kept and called after the flush. For setups where it cannot attach, pass listener=ConsumerRebalanceListener.from_context(broker.context) to the subscriber. When Kafka revokes a partition, the listener calls committer.commit_all() to flush pending offsets before the partition is reassigned. The flush waits for in-flight tasks up to rebalance_flush_timeout_sec (default 10 s) so a slow handler cannot stall the rebalance past max.poll.interval.ms; on timeout, the remaining in-flight messages are redelivered after reassignment (at-least-once). A future optimization may scope the wait to only the revoked partitions.

Shutdown

stop_concurrent_processing cancels every in-flight asyncio task, then awaits committer.close(). The committer treats cancelled tasks as a hard offset boundary: the cancelled offset and the ones after it stay uncommitted and get redelivered on restart. Total wall-clock is sub-second in normal conditions and bounded by shutdown_timeout_sec only as a safety net for stuck network commits. If your handlers do non-idempotent work that is expensive to repeat, wrap them in try/finally so cleanup runs on CancelledError.

Limitations

FastStream control signals from a middleware registered after this one

A middleware you register after KafkaConcurrentProcessingMiddleware runs inside the coroutine this library dispatches as a background task, so a FastStream control signal it raises never reaches FastStream. Every such signal is absorbed and the message's offset is committed; what differs is whether the library could act on it.

raised by an inner middleware effect logged
AckMessage offset commits; this is the ack DEBUG
RejectMessage offset commits (for Kafka, reject() is ack()) DEBUG
SkipMessage offset commits, processing moves on DEBUG
NackMessage not honoured: offset commits instead of being redelivered ERROR
StopConsume not honoured: the subscriber keeps consuming ERROR
StopApplication not honoured: the application keeps running ERROR

The three marked not honoured do not work from a concurrently dispatched handler. They log an ERROR naming the signal. If you depend on any of them, raise it from a middleware registered before KafkaConcurrentProcessingMiddleware (which runs outside the dispatched task), or from outside the message-processing path entirely.

Rationale and the rejected alternatives: ADR-0002.

Calling msg.ack() / msg.nack() / msg.reject() directly

On the concurrent path these calls raise RuntimeError. Offset control belongs to KafkaBatchCommitter, and reaching around it silently loses data.

KafkaAckableMessage.ack() issues a bare consumer.commit() with no offsets, committing the consumer's current fetch position, which is past every in-flight task on every assigned partition, so those messages are never processed and never redelivered. reject() is an ack for Kafka and carries the same hazard under an opposite-sounding name. nack() issues consumer.seek(...), rewinding the partition underneath tasks already processing it.

There is no supported way to request redelivery under concurrent processing: the offset commits even when your handler raises. See ADR-0002.

Subscribers that pass through (a FakeConsumer under TestKafkaBroker, or any non-MANUAL ack policy) are unaffected, because this library is not managing their offsets. The guards are installed only on the concurrent dispatch path, so a passed-through subscriber never sees them.

Reaching through the message to the raw consumer, as in msg.consumer.commit() or msg.consumer.seek(...), is still unguarded. The consumer is one shared object across every message and partition, so it cannot be guarded per message. Do not do it.

Other

  • Batch subscribers (batch=True) are unsupported. A batch=True subscriber declaring AckPolicy.MANUAL is rejected with an explicit RuntimeError, because the concurrent path is one message → one task → one offset. A batch=True subscriber on any other ack policy passes through, since the middleware does not manage it at all.
  • ack_policy=AckPolicy.MANUAL is required on subscribers you want processed concurrently. Every other policy passes through untouched, exactly as if the middleware were not registered, which makes a single broker-level add_middleware call safe across a mix of subscribers. ACK_FIRST leaves its offsets to aiokafka's enable_auto_commit; ACK, REJECT_ON_ERROR and NACK_ON_ERROR are acknowledged by FastStream's own AcknowledgementMiddleware as soon as this middleware returns. That ack is safe on the pass-through path: each FastStream subscriber builds its own AIOKafkaConsumer, so it touches only that subscriber's partitions and cannot commit past another subscriber's in-flight work, and with no background task "consumed" and "processed" are the same moment.

Requirements

  • Python >= 3.11
  • faststream[kafka]

📦 PyPI

📝 License

Part of modern-python

Browse the full list of templates and libraries in modern-python; the org profile has the categorized index.

About

Concurrent message-processing middleware for FastStream + aiokafka

Topics

Resources

Code of conduct

Contributing

Security policy

Stars

3 stars

Watchers

0 watching

Forks

Releases

Sponsor this project

Contributors

Languages