From b88a3c7aafc6a691de3f0f7252f37e4c85fdd07f Mon Sep 17 00:00:00 2001 From: Artur Shiriev Date: Sat, 3 Oct 2026 17:26:58 +0300 Subject: [PATCH] docs: plain-prose pass over README --- README.md | 105 +++++++++++++++++++++++++++++------------------------- 1 file changed, 57 insertions(+), 48 deletions(-) diff --git a/README.md b/README.md index d0a55fa..2bc9512 100644 --- a/README.md +++ b/README.md @@ -20,7 +20,7 @@ Concurrent message processing middleware for [FastStream](https://faststream.ag2.ai/) 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. +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 @@ -39,14 +39,14 @@ By default FastStream processes Kafka messages sequentially — one message at a pip install faststream-concurrent-aiokafka ``` -## Quick Start +## Quick start -`ack_policy=AckPolicy.MANUAL` is required for concurrent processing; other policies pass through untouched. +`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` — are passed through without concurrent processing, -behaving exactly as they would if this middleware were not registered. That keeps a single broker-level registration safe across a mix of subscribers. +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. -> **`AsgiFastStream` note**: its lifespan receives an app-level `ContextRepo` separate from `broker.context`. Pass `broker.context` explicitly instead of the injected argument. +> With `AsgiFastStream`, the lifespan receives an app-level `ContextRepo` separate from `broker.context`. Pass `broker.context` explicitly instead of the injected argument. ```python from contextlib import asynccontextmanager @@ -92,7 +92,7 @@ async def handle(msg: str) -> None: ... async def handle_other(msg: str) -> None: ... ``` -## Core Concepts +## Core concepts `KafkaConcurrentHandler` and `KafkaBatchCommitter` are internal: they are not exported from the package and are described here only to explain the behavior. @@ -102,20 +102,20 @@ A FastStream `BaseMiddleware` subclass. Add it to your broker to enable concurre ### KafkaConcurrentHandler -The processing engine. Manages: +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](#faststream-control-signals-from-a-middleware-registered-after-this-one) 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](#how-it-works) +- 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](#rebalance-handling) -This library does **not** install signal handlers — shutdown is driven by your lifespan / process manager calling `stop_concurrent_processing`. +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 treated as 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. +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 +## API reference ### `initialize_concurrent_processing(context, ...)` @@ -133,11 +133,11 @@ Create and start the concurrent processing handler; store it in FastStream's con Returns the `KafkaConcurrentHandler` instance. -> **Tuning `max_uncommitted_tasks`:** each uncommitted entry holds only commit metadata — a task reference, its `TopicPartition`, offset, and consumer reference — not the message payload, so the default 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 restore unbounded buffering. +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. +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)` @@ -145,15 +145,15 @@ Returns `True` if the `KafkaConcurrentHandler` stored in `context` is running an ### `KafkaConcurrentProcessingMiddleware` -FastStream middleware class. Register it via `broker.add_middleware(...)`. See Quick Start for usage examples. +FastStream middleware class. Register it via `broker.add_middleware(...)`. See [Quick start](#quick-start) for usage examples. -> **Must be outermost.** `consume_scope` fires the handler as a background task and returns `None` immediately. Any middleware that wraps it on the outside will see that premature return and misfire — wrong timing, early cleanup, or missed exceptions. Middlewares added after it (i.e. inner in the chain) run correctly inside the background task. +> 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. +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. -**Fix**: 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. +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. ```python broker = KafkaBroker(...) @@ -161,23 +161,33 @@ broker.add_middleware(KafkaConcurrentProcessingMiddleware) # registered first modern_di_faststream.setup_di(app, container=container) # registered after → inner to KCM ``` -## How It Works +## How it works -1. **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`. +### Message dispatch -2. **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. +On each incoming message, `consume_scope` calls `handle_task()`, which acquires a semaphore slot then fires the handler coroutine as a background `asyncio.Task`. -3. **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). +### Concurrency control -4. **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. +The semaphore blocks new tasks when `concurrency_limit` is reached. The slot is released via a done-callback when the task finishes or fails. -5. **Shutdown**: `stop_concurrent_processing` cancels every in-flight asyncio task, then awaits `committer.close()`. The committer treats cancelled tasks as a hard offset boundary — cancelled-and-after offsets 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`. +### 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 +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 @@ -185,16 +195,16 @@ 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 | +| `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 | +| `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** +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. @@ -203,13 +213,12 @@ Rationale and the rejected alternatives: ### Calling `msg.ack()` / `msg.nack()` / `msg.reject()` directly -**These raise `RuntimeError` on the concurrent path.** Offset control belongs to -`KafkaBatchCommitter`; reaching around it silently loses data, so the middleware -raises. +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* — past every in-flight task on -every assigned partition — so those messages are never processed and never +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. @@ -218,31 +227,31 @@ There is no supported way to request redelivery under concurrent processing: the offset commits even when your handler raises. See [ADR-0002](https://github.com/modern-python/faststream-concurrent-aiokafka/blob/main/docs/adr/0002-control-signals-not-honoured.md). -Subscribers that pass through — a `FakeConsumer` under `TestKafkaBroker`, or any -non-`MANUAL` ack policy — are unaffected, because this library is not managing +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. -**Still unguarded:** reaching through the message to the raw consumer, as in -`msg.consumer.commit()` or `msg.consumer.seek(...)`. The consumer is one shared +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`. The concurrent path is one message → one task → one offset. A - `batch=True` subscriber on any other ack policy simply passes through, since +- 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 +- `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 — that is what makes a single broker-level + 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 + 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. @@ -259,4 +268,4 @@ Do not do it. ## Part of `modern-python` Browse the full list of templates and libraries in -[`modern-python`](https://github.com/modern-python) — see the org profile for the categorized index. +[`modern-python`](https://github.com/modern-python); the org profile has the categorized index.