How should producer creation handle backlog-quota rejection? #26682
Unanswered
jaswanthnaidumalasanimalasani-code
asked this question in
Q&A
Replies: 0 comments
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Uh oh!
There was an error while loading. Please reload this page.
We publish through an HTTP endpoint that creates producers on demand, one per topic, and caches them with a 30s TTL.
The relevant setup is:
PulsarClient, shared across all topics and callers.Sharedaccess mode, batching and compression are disabled, andmaxPendingMessages(500)withblockIfQueueFull(false).The 30s TTL means an intermittently used topic can leave the cache and trigger producer creation again on a later request.
When a topic reaches its backlog quota with
producer_exception, the broker refuses new producers. That behavior is expected.What I am trying to understand is how the caller is expected to handle this refusal. The client appears to treat the quota exception as retriable and continues producer creation instead of surfacing the exception immediately.
The behavior below is not specific to an old release. I checked
masterand the current release branches, and the relevant code is the same on all of them.Observed behavior
With:
the retry sequence is approximately:
mandatoryStop.Attempt 10 is admitted just before the 30s deadline, but the subsequent backoff is ~51s. As a result, producer creation fails after ~81s even though
operationTimeoutis 30s.The exact timings vary slightly because of backoff jitter.
What I found in the client
isRetriableErrorappears to be a denylist. Neither of these exceptions is explicitly excluded:ProducerBlockedQuotaExceededExceptionProducerBlockedQuotaExceededErrorProducerImplhas a specific branch for the quota case:Topic backlog quota exceeded. Throwing Exception on producer.and calls
failPendingMessages, but it does not completeproducerCreatedFutureexceptionally.The
TopicTerminatedExceptionandProducerFencedExceptioncases immediately below it do:The quota case instead reaches the generic retriable path and schedules another reconnect.
Also,
operationTimeoutappears to be checked when deciding whether another attempt can start. Once an attempt has been admitted, the subsequent backoff is scheduled without applying the remaining deadline to that delay. This is what allows the attempt at ~29.9s to result in another ~51s wait.mandatoryStoponly truncates the backoff crossing the stop point; subsequent backoff continues normally.What we tried
We originally created producers synchronously, so HTTP request threads could stay blocked in
ProducerBuilderImpl.create()for the duration of the retry sequence. We have tried two ways to avoid that, and neither feels like the intended answer.1. Create asynchronously, with our own deadline. We switched to
createAsync()and wrapped it inorTimeout(...)so the request thread is never held. This solves the blocking, but not the underlying behavior: creation for a quota-exhausted topic stays pending until our deadline fires, and requests waiting on it accumulate along with their payloads. The request-thread pool had been an implicit bound on how many could be outstanding, and removing the blocking removed that bound too.2. Keep creation synchronous, but reduce
operationTimeoutandmaxBackoffInterval. With a lower backoff ceiling, every sleep is short enough that the typed quota exception surfaces within our timeout, so we can recognise it and stop attempting that topic for a while. This works. But it only works because one value happens to fit inside the other, across two settings that are otherwise unrelated, on aPulsarClientshared by every topic we publish to. Nothing enforces that relationship, and changing either value silently changes whether the exception is observable.Questions
Is it intentional that
ProducerBlockedQuotaExceededException/ProducerBlockedQuotaExceededErrorare treated as retriable?In particular, is this because the same retry path is also used for
producer_request_hold, where waiting for the quota to clear is expected?Is there a supported way to bound the total producer creation time?
operationTimeoutappears to control retry admission rather than the completecreateAsync()operation. Is wrapping it in an application-level deadline, as in approach 1, the expected way to do this, or is there a client setting for it that I have missed?Is reducing
maxBackoffIntervalthe expected way to make the quota exception observable within a desired timeout?For an HTTP producer service, is sharing one
PulsarClientacross topics/tenants the recommended approach, or is there a recommended client scope for isolating this kind of retry behavior?Is there a recommended way to determine that a topic is already at its backlog quota before attempting producer creation?
We build the client directly using
PulsarClient.builder()rather than through the WebSocket proxy.Related issues I found are #26373 and #26343, although I don't see either covering this producer-creation/retriability behavior.
I can provide the relevant thread dumps or code paths if useful.
All reactions