diff --git a/README.md b/README.md index 2bc9512..d7d6943 100644 --- a/README.md +++ b/README.md @@ -177,7 +177,7 @@ Each dispatched task is paired with its Kafka offset and consumer reference and ### 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. +`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). ### Shutdown