diff --git a/README.md b/README.md index 818f186..d0a55fa 100644 --- a/README.md +++ b/README.md @@ -29,7 +29,7 @@ By default FastStream processes Kafka messages sequentially — one message at a - 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) -- Signal handling owned by your lifespan / process manager — this lib does not register SIGTERM/SIGINT handlers +- 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` @@ -41,7 +41,7 @@ pip install faststream-concurrent-aiokafka ## Quick Start -`ack_policy=AckPolicy.MANUAL` is **required** on every concurrent subscriber — the middleware enforces this at runtime. +`ack_policy=AckPolicy.MANUAL` is required for concurrent processing; other policies pass through untouched. 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. @@ -70,7 +70,7 @@ 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 completed tasks + commit_batch_size=100, # commit after this many pending tasks commit_batch_timeout_sec=5.0, # or after this many seconds ) try: @@ -94,6 +94,8 @@ 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`. @@ -123,7 +125,7 @@ Create and start the concurrent processing handler; store it in FastStream's con |---|---|---| | `context` | required | FastStream `ContextRepo` instance | | `concurrency_limit` | `10` | Max concurrent asyncio tasks (minimum: 1) | -| `commit_batch_size` | `10` | Max messages per commit batch | +| `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. | @@ -167,9 +169,9 @@ modern_di_faststream.setup_di(app, container=container) # registered after → 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). -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 `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. +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. -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. +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`. ## Limitations @@ -190,10 +192,8 @@ library could act on it. | `StopConsume` | **not honoured** — the subscriber keeps consuming | ERROR | | `StopApplication` | **not honoured** — the application keeps running | ERROR | -The three marked *not honoured* have never worked from a concurrently dispatched -handler. Before 0.6.4 they failed silently — and `StopApplication`, which -subclasses `SystemExit`, tore down the event loop outright, losing every -in-flight offset. They now log a loud ERROR naming the signal instead. If you +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. @@ -205,7 +205,7 @@ Rationale and the rejected alternatives: **These raise `RuntimeError` on the concurrent path.** Offset control belongs to `KafkaBatchCommitter`; reaching around it silently loses data, so the middleware -refuses the call rather than letting it through. +raises. `KafkaAckableMessage.ack()` issues a bare `consumer.commit()` with no offsets, committing the consumer's *current fetch position* — past every in-flight task on @@ -245,21 +245,7 @@ Do not do it. 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. It would *not* be safe if such - a subscriber were dispatched, which is precisely why it is not. - -## Migration from < 0.x - -Previously, `stop_concurrent_processing` waited up to `2 × shutdown_timeout_sec` for in-flight tasks to drain to completion. The new behavior cancels them immediately. The at-least-once contract is unchanged — uncommitted offsets are redelivered on restart, the same way they always were when the handler crashed mid-task. - -| What changed | Old | New | -|---|---|---| -| In-flight tasks on stop | drained to completion | **cancelled** | -| `KafkaConcurrentHandler.wait_for_subtasks()` | public method | removed | -| `shutdown_timeout_sec` | applied separately to handler and committer | applied to committer only | -| Signal handler installation | installed automatically | removed — own them via your lifespan / process manager | - -If your handlers do non-idempotent work that's expensive to repeat, ensure your handlers are wrapped in `try/finally` so cleanup runs on `CancelledError`, or pin to the previous version of this library. To trigger shutdown on SIGTERM/SIGINT, your lifespan or main entry point must catch the signal and call `stop_concurrent_processing(broker.context)` — under uvicorn / AsgiFastStream this happens automatically through the lifespan `finally` block. + "consumed" and "processed" are the same moment. ## Requirements