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.
- 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
pip install faststream-concurrent-aiokafkaack_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-levelContextReposeparate frombroker.context. Passbroker.contextexplicitly 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: ...KafkaConcurrentHandler and KafkaBatchCommitter are internal: they are not exported from the package and are described here only to explain the behavior.
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.
The processing engine. It manages:
- An
asyncio.Semaphoreto enforceconcurrency_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
KafkaBatchCommitterfor offset commits - A
ConsumerRebalanceListeneron every concurrent subscriber that flushes pending commits when partitions are revoked.initialize_concurrent_processingattaches 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.
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.
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.
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).
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.
FastStream middleware class. Register it via broker.add_middleware(...). See Quick start for usage examples.
This middleware must be the outermost one.
consume_scopefires the handler as a background task and returnsNoneimmediately, 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 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 KCMOn each incoming message, consume_scope calls handle_task(), which acquires a semaphore slot then fires the handler coroutine as a background asyncio.Task.
The semaphore blocks new tasks when concurrency_limit is reached. The slot is released via a done-callback when the task finishes or fails.
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).
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.
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.
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.
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.
- Batch subscribers (
batch=True) are unsupported. Abatch=Truesubscriber declaringAckPolicy.MANUALis rejected with an explicitRuntimeError, because the concurrent path is one message → one task → one offset. Abatch=Truesubscriber on any other ack policy passes through, since the middleware does not manage it at all. ack_policy=AckPolicy.MANUALis 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-leveladd_middlewarecall safe across a mix of subscribers.ACK_FIRSTleaves its offsets to aiokafka'senable_auto_commit;ACK,REJECT_ON_ERRORandNACK_ON_ERRORare acknowledged by FastStream's ownAcknowledgementMiddlewareas soon as this middleware returns. That ack is safe on the pass-through path: each FastStream subscriber builds its ownAIOKafkaConsumer, 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.
- Python >= 3.11
faststream[kafka]
📦 PyPI
📝 License
Browse the full list of templates and libraries in
modern-python; the org profile has the categorized index.