Conversation
Lancetnik
left a comment
There was a problem hiding this comment.
Thanks for picking this up, and for the reproduction in #2836 — the diagnosis there is accurate and the regression test here is a genuine one (it does fail without the patch). The problem is real and worth fixing. My concern is the direction rather than the details, so before another round on this diff, I'd like to lay out what I think it costs.
The other half of #2836 is untouched
#2836 reports two failures: the cancellation cascade, and 100 000 InvalidStateError from orphaned delivery callbacks. This patch removes one trigger of the cascade, but the orphaned-callback path is still live. The root is in ack_callback:
https://github.com/ag2ai/faststream/blob/main/faststream/confluent/helpers/client.py#L115-L122
It calls set_result / set_exception on a future that a sibling failure may already have cancelled. Any error inside send — a delivery error on one message, for instance — still cancels the task group and reproduces the same storm; BufferError was only the most reachable trigger. One line fixes it for every path:
def ack_callback(err: Any, msg: Message | None) -> None:
if result_future.done(): # cancelled by a sibling failure
return
...That is worth doing regardless of which direction the rest of this takes.
poll() in the retry loop competes for the thread pool
call_or_await(self.producer.poll, 1.0) resolves to anyio.to_thread.run_sync, whose default limiter holds 40 tokens:
>>> anyio.to_thread.current_default_thread_limiter().total_tokens
40So each retrying message occupies one of 40 threads for a full second. A batch that overshoots the queue by 100k messages puts 100k coroutines into that loop, and they drain the pool that AsyncConfluentConsumer.getone / getmany and the producer's own _poll_loop also use. The publish path starves the consume path of an unrelated subscriber in the same app.
The redundancy makes it avoidable: _poll_loop already polls every 100 ms and already serves delivery reports:
https://github.com/ag2ai/faststream/blob/main/faststream/confluent/helpers/client.py#L73-L76
An await anyio.sleep(...) would let the existing loop drain the queue and hold no thread at all.
while True has no exit
There is no attempt limit, no deadline, and no log line. If the queue stops draining — broker unreachable, with message.timeout.ms at its 5-minute default — publish() blocks with no way for the caller to tell a stall from normal operation. Compare faststream.kafka, which fails fast with a position. Whatever the retry ends up looking like, it needs a bound and a message on the first BufferError.
Single publish() changes behaviour too
The patch sits in send(), not send_batch(), so it also covers non-batch publishing. BufferError used to surface to the caller immediately; now a single publish() blocks indefinitely instead. That is outside the stated scope of the change and isn't mentioned in the description.
Minor: _BUFFER_FULL_POLL_TIMEOUT = 1.0 is unrelated to both ConfluentFastConfig and the 0.1 s step of _poll_loop.
What I'd suggest instead
Chunking is proactive where retrying is reactive — it keeps the queue from overflowing rather than recovering after it has:
- Add the
ack_callbackguard above. One line, independent of everything else. - Chunk
send_batchbyself.config.get("queue.buffering.max.messages", 100000), withflush()between chunks — roughly what you sketched in #2836. - Keep a retry on
BufferErroras a safety net forqueue.buffering.max.kbytes, which chunking by count cannot cover — but bounded, logged, and sleeping rather than polling in a thread.
(1) and (2) alone close both halves of #2836 and are a smaller diff than this one. If you would rather land (1) on its own first, that is a clean standalone PR and I will merge it quickly.
For context on where this path leads longer term: confluent_kafka.Producer.produce_batch would replace this whole code path with one call that reports partial failure natively, no task group involved. It is blocked on missing header support upstream — tracked in #2993 — so it is not an option today, but it is the shape we are aiming at, and chunking is compatible with it.
|
Thanks for the detailed review, this was really useful. Reworked it along the lines you laid out: the done() guard in ack_callback, send_batch now chunks by queue.buffering.max.messages with a flush between chunks, and the BufferError retry is kept only as a safety net for the kbytes limit in the batch path, bounded by message.timeout.ms, logged on the first hit, and just sleeping so the existing poll loop drains the queue instead of tying up a thread. Single publish() raises BufferError immediately again, so no semantic change outside the batch path. Adapted the regression test and added coverage for the chunking, the fail-fast path, and the callback guard. Also removed the stray timeout constant. |
There was a problem hiding this comment.
The rework is right on every point from the last round — the done() guard, chunking with a flush between chunks, sleeping instead of polling in a thread, and single publish() back to failing fast. Thank you for turning it around so quickly.
One change before this lands: retry_on_buffer_error should be exposed on the public publish_batch, defaulting to False.
As it stands no caller can reach it, so batch publishing has exactly one behaviour — wait up to message.timeout.ms, 5 minutes by default. After chunking, an overflow means the queue.buffering.max.kbytes limit was hit, which is rare and worth surfacing; a caller who would rather re-split the batch or fail loudly instead sees a stall with nothing to act on. It also splits the two Kafka backends: faststream.kafka raises BatchBufferOverflowException and hands the batch back, so the same condition gives a signal on one backend and a five-minute wait on the other.
What we want: publish_batch fails fast on BufferError by default, and a caller who prefers waiting for the queue to drain opts into it explicitly.
|
Fair point, agreed the batch path shouldn't silently sit on a full queue by default. Exposed retry_on_buffer_error on publish_batch (broker and batch publisher), defaulting to False, so BufferError now surfaces immediately like the kafka backend and the drain-and-retry is an explicit opt-in threaded down to send_batch. Adjusted the tests accordingly and added one covering the fail-fast default. |
The branch carried a verbatim copy of ag2ai#3012, which reached main as a squash merge. With no shared ancestry the 3-way merge added the same three MQTTBroker(...) calls a second time instead of flagging a conflict. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
There was a problem hiding this comment.
Everything from the last round is closed — I checked each point, including that retry_on_buffer_error reaches AsyncConfluentProducer.send_batch from both KafkaBroker.publish_batch and BatchPublisher.publish with False as the default. Six things left.
1. Put chunking behind retry_on_buffer_error as well, so the default path stays exactly what it is today.
This one is on me — last round I asked for chunking as the proactive fix and you implemented precisely that. But it is currently unconditional, so every caller who never touches the new flag gets new behaviour once a batch crosses queue.buffering.max.messages. Three ways it differs from main: no_confirm=True stops meaning what it says, because we now wait for delivery between chunks (and only between chunks — the last one is not flushed, so the batch is half fire-and-forget); flush() drains the entire producer queue, which is shared across all topics and partitions and across every publisher on the broker, so a large batch blocks on unrelated in-flight messages; and it is not cancellable, since await self.flush() goes through run_in_threadpool without cancellable=True and calls producer.flush() with no timeout, so an unreachable broker holds a thread until message.timeout.ms with no way to interrupt.
Keep the ack_callback done() guard unconditional though — that is a plain bug fix, and gating it would leave the second half of #2836 open for anyone who does not opt in.
2. Guard the chunk size against 0.
ConfluentFastConfig(config={"queue.buffering.max.messages": 0})
# ValueError: range() arg 3 must not be zero0 is a valid librdkafka value meaning "no limit" — Producer({"queue.buffering.max.messages": 0}) constructs fine on librdkafka 2.14.0. That config works today; with this patch every publish_batch raises.
3. Do not read message.timeout.ms=0 as a plain number either.
It means "no delivery timeout" in librdkafka, but here it becomes deadline = now, so the retry gives up after the first 100 ms sleep — the opposite of what the setting asks for. Same root cause as (2).
4. Log the BufferError warning once per send_batch, not once per message.
500 overflowing messages produce 500 warning lines; at a 100 000-message chunk that is 100 000 of them, which is the log storm #2836 opened with.
5. Document retry_on_buffer_error.
It appears nowhere under docs/. A public keyword whose effect is "this call may block for up to message.timeout.ms, five minutes by default" should not be discoverable only from the signature. Please cover what triggers it — queue.buffering.max.kbytes, which count-based chunking cannot catch — and what bounds it.
6. Test through the public API wherever a case can be expressed there.
All five tests drive AsyncConfluentProducer.send_batch directly with a fake Producer. That is the right level for the chunk-count and done()-guard assertions, but it leaves the thing this whole round was about — that retry_on_buffer_error is reachable at all — with no coverage. I confirmed the plumbing by hand, so a regression there would ship silently. Please add cases at broker.publish_batch(...) and broker.publisher(batch=True).publish(...) asserting both the opt-in and the default.
|
Thanks for another really thorough pass, all six are in. Chunking (and the between-chunk wait) now only kicks in with retry_on_buffer_error=True, so the default path is byte-identical to main again: no chunking, no flush, no_confirm means what it says. The done() guard stays unconditional as you asked. The wait between chunks no longer goes through flush() at all, it polls the local queue length with short async sleeps, so nothing holds a threadpool thread and it cancels cleanly. Both zero-value configs are handled (queue.buffering.max.messages=0 disables chunking, message.timeout.ms=0 retries with no deadline), the warning fires once per send_batch, the flag is documented on the batch publisher page, and there are now public-API tests through broker.publish_batch and publisher(batch=True).publish covering both the opt-in and the fail-fast default. |
Description
This supersedes #2941, which was closed while my CLA signature was still missing. The CLA is signed now, and the change still applies cleanly on current main.
Fixes #2836.
When a batch pushes librdkafka's local produce queue past
queue.buffering.max.messages,produce()raisesBufferError("Local: Queue full"). That error was bubbling up out of a singlesendtask and cancelling the wholesend_batchtask group, so one overflowing message took down every sibling message in the batch.This catches
BufferErrorin the produce path, serves delivery reports withpoll()to drain the queue, then retries that same message. It's the usual confluent-kafka backpressure pattern, so the rest of the batch goes through instead of being cancelled.Type of change
Checklist
Added
tests/brokers/confluent/test_buffer_full.py, which fakes a full queue on the first produce call and checks the batch still completes (it fails without the fix). That test and the existing confluent unit tests pass locally.