Skip to content

fix(confluent): retry on BufferError so a full queue doesn't cancel the whole batch - #2986

Open
cycsmail wants to merge 8 commits into
ag2ai:mainfrom
cycsmail:fix-confluent-buffer-full-batch-2836
Open

cycsmail wants to merge 8 commits into
ag2ai:mainfrom
cycsmail:fix-confluent-buffer-full-batch-2836

Conversation

@cycsmail

@cycsmail cycsmail commented Aug 4, 2026

Copy link
Copy Markdown

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() raises BufferError ("Local: Queue full"). That error was bubbling up out of a single send task and cancelling the whole send_batch task group, so one overflowing message took down every sibling message in the batch.

This catches BufferError in the produce path, serves delivery reports with poll() 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

  • Bug fix (a non-breaking change that resolves an issue)

Checklist

  • I have conducted a self-review of my own code
  • I have added tests to validate the effectiveness of my fix

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.

@cycsmail
cycsmail requested a review from Lancetnik as a code owner August 4, 2026 03:09
@github-actions github-actions Bot added the Confluent Issues related to `faststream.confluent` module label Aug 4, 2026
@ApusBerliozi ApusBerliozi added the bug Something isn't working label Aug 4, 2026

@Lancetnik Lancetnik left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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
40

So 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:

  1. Add the ack_callback guard above. One line, independent of everything else.
  2. Chunk send_batch by self.config.get("queue.buffering.max.messages", 100000), with flush() between chunks — roughly what you sketched in #2836.
  3. Keep a retry on BufferError as a safety net for queue.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.

@cycsmail
cycsmail requested a review from Sehat1137 as a code owner August 6, 2026 09:18
@cycsmail

cycsmail commented Aug 6, 2026

Copy link
Copy Markdown
Author

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.

@github-actions github-actions Bot added the dependencies Pull requests that update a dependency file label Aug 6, 2026

@Lancetnik Lancetnik left a comment •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@cycsmail

Copy link
Copy Markdown
Author

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.

@github-actions github-actions Bot added documentation Improvements or additions to documentation Redis Issues related to `faststream.redis` module and Redis features MQTT Issues related to `faststream.mqtt` module labels Aug 14, 2026
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>

@Lancetnik Lancetnik left a comment •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 zero

0 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.

@cycsmail
cycsmail requested a review from powersemmi as a code owner September 1, 2026 05:20
@cycsmail

cycsmail commented Sep 1, 2026

Copy link
Copy Markdown
Author

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.

@github-actions github-actions Bot added github_actions Pull requests that update GitHub Actions code AioKafka Issues related to `faststream.kafka` module NATS Issues related to `faststream.nats` module and NATS broker features AsyncAPI Issues related to AsyncAPI specification generation labels Sep 1, 2026

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

AioKafka Issues related to `faststream.kafka` module AsyncAPI Issues related to AsyncAPI specification generation bug Something isn't working Confluent Issues related to `faststream.confluent` module dependencies Pull requests that update a dependency file documentation Improvements or additions to documentation github_actions Pull requests that update GitHub Actions code MQTT Issues related to `faststream.mqtt` module NATS Issues related to `faststream.nats` module and NATS broker features Redis Issues related to `faststream.redis` module and Redis features

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Bug: confluent publish_batch — one message over queue.buffering.max.messages cancels entire batch via TaskGroup cascade

3 participants