Repository navigation
[fix][broker] Fix persistent throughput degradation caused by permit loss during frequent reconnects on Shared subscriptions - #26289
Conversation
lhotari
left a comment
There was a problem hiding this comment.
Reviewed this in depth and verified it against a local build. The core fix is correct - absent integer overflow (see the inline note on the Math.max clamp) I could not find a case where it makes accounting worse. Everything below is non-blocking, on top of @Denovo1998's review, which I think is asking the right questions.
What I verified
The diagnosis is right, and the failure is stronger than a narrow race. ServerCnx.handleFlow runs on the connection event loop and defers the dispatcher update to the broker executor, while handleCloseConsumer -> Consumer.close -> PersistentSubscription.removeConsumer -> dispatcher.removeConsumer all run synchronously on that same event loop. Any close shortly after a Flow hits this while the queued task is still pending - exactly the reconnect workload in #26288. The unload path is covered too, since disconnectAllConsumers holds the dispatcher monitor across consumer.disconnect() -> removal.
The symptom matches the issue. readMoreEntries floors the read at Math.max(totalAvailablePermits, getFirstAvailableConsumerPermits()), so a negative total does not stall the subscription - with several consumers connected it shrinks each read toward a single consumer's permits instead of their sum. That is the "30-40% slower but still progressing" behaviour in #26288 rather than a hard stop.
The invariant restored is exact: totalAvailablePermits == sum(messagePermits - pendingDispatcherFlowPermits) over connected consumers. Checked at every mutation site: both internalConsumerFlow branches, Consumer.sendMessages against the matching TOTAL_AVAILABLE_PERMITS_UPDATER decrements in both dispatchers and the sticky-key one, all three removeConsumer branches, clearComponentsAfterRemovedAllConsumers, and the blocked-permits path (PERMITS_RECEIVED_WHILE_CONSUMER_BLOCKED correctly sits in neither side). The accounting is also commutative: BrokerService.executor() round-robins Flow tasks across threads so they can complete out of order, but each task completes exactly what it added, so every interleaving converges.
Threading is sound. flowPermitAccountingLock is a leaf lock - nothing is acquired under it and no callback or I/O runs inside - so there is no ordering relationship with the dispatcher monitor to get wrong, and the new critical sections are two field updates.
On @Denovo1998's per-Flow overhead concern: it is a per-Consumer monitor, contended only between that consumer's event loop and the broker executor. flowPermits already calls System.currentTimeMillis() and submits a task to an executor in the same method, so the uncontended monitor is small next to what is already there. I do not think overhead alone justifies restructuring, though the correctness argument for scoping still stands.
Executor rejection is correct by construction - I went looking for a leak and there isn't one. If consumerFlow is ever rejected at shutdown the flow legitimately stays pending, and removal then subtracts exactly the permits that were applied.
Tests pin the fix. On master with only SharedDispatcherPermitAccountingTest added, all three methods fail with expected: 10 but was: -990; on this branch all 5 invocations pass in 3.3s. :pulsar-broker:checkstyleMain :pulsar-broker:checkstyleTest pass.
On the open review points
- The counter-scoping question (Consumer.java:931) is real, and I have added inline the consequence that I think makes it worth acting on: the accessor becomes an active trap for whoever extends this fix to the other dispatchers.
- The legitimately-negative balance (Consumer.java:1000) is real and load-bearing. I have posted the mechanism and a test shape inline. Note it needs no batching or special configuration - it follows from dispatch being sized off
getAvailablePermits(), which includes permits the dispatcher has not counted yet. - On Key_Shared coverage, agreed - the three implementations are listed inline. Their
removeConsumeroverrides wrapsuper.removeConsumerwith selector and draining-hash work, so a Key_Shared variant of the race test is worth having even though the permit arithmetic itself is inherited unchanged.
Out of scope, but worth not losing
The same defect class survives on non-persistent Shared subscriptions: NonPersistentDispatcherMultipleConsumers.removeConsumer still subtracts the full getAvailablePermits(), and its consumerFlow drops updates from consumers already out of consumerSet. The window is narrower because that consumerFlow is synchronous, but not empty - NonPersistentTopic.onPoliciesUpdate -> Consumer.checkPermissionsAsync -> disconnect() -> close() runs removal on the authorization future's completion thread rather than the connection event loop, and disconnectAllConsumers holds the dispatcher monitor across the whole teardown. There the consequence is worse than a slowdown: sendMessages drops entries when the total is not positive. Full teardowns self-heal, since the total is reset to 0 once the consumer list empties, so lasting drift needs a partial removal that leaves survivors.
I would keep that out of this PR - it needs the scoping question settled first - but it is worth a follow-up issue.
Backport note
The release/4.2.5 and release/4.0.14 labels are on this PR: the diff touches slog-style logging (log.debug().attr(...)), and branch-4.2 / branch-4.0 are Maven + slf4j, so those cherry-picks will need the usual logging adaptation rather than a clean pick.
|
@lhotari Thanks for the detailed review and the follow-up suggestions. I have pushed I agree that the non-persistent Shared case should be handled separately. I will first check whether it is related to #24018 before opening a new issue. I have also noted the backport point. Since the release labels are already on the PR, no further action is needed here, and I can help with branch-specific logging adaptation if needed. |
lhotari
left a comment
There was a problem hiding this comment.
The follow-up commit addresses all five review threads. The pending counter now preserves signed wrap, tracking is restricted to persistent Shared/Key_Shared consumers, the accessor documents its precondition and signed invariant, and the tests cover wrapped and negative balances, both dispatcher variants, queued-Flow ordering, and production lock order. I found no new issue in this review-feedback-only delta.
We discussed this with @merlimat and took a look together. In |
|
@lhotari Thanks for the suggestion and for looking into this with @merlimat. Updated in
|
The inline Flow path can wait for the persistent Shared dispatcher monitor on the connection EventLoop while dispatch or filter work owns the monitor. Route Flow accounting through the dispatcher's selected dispatchMessagesThread in both modern and classic implementations. Complete pending accounting and update total permits under the monitor, then trigger the read from the same lane. Leave rejected tasks pending so removal still excludes unapplied permits. Update the race tests for queued Flow processing and add deterministic coverage for EventLoop progress, lane affinity, removal, and shutdown rejection across Shared and Key_Shared.
|
Added follow-up commit With The follow-up changes the threading model to:
No new executor or thread pool is introduced. Validation:
Receiver queue size 1 was used as an intentional boundary stress case, generating approximately one Flow per delivered message. Throughput remained stable, receive-gap p99 differed by less than approximately 0.5%, and all correctness/progress checks passed. In the q10 control, the measured CPU difference was approximately 3%. The result is that Flow processing no longer makes the connection EventLoop wait for the dispatcher monitor, without an observed correctness, progress, or sub-saturation throughput/latency regression. |
A queued Flow task can outlive its Consumer and run after a replacement with the same protocol identity has joined. Consumer equality can then make the old task appear connected and add stale permits. Keep the existing ObjectSet field while retaining its ObjectHashSet backing instance for an O(1) identity check. Use that check for modern and classic queued Flow processing and cover same-identity replacement across Shared and Key_Shared.
Denovo1998
left a comment
There was a problem hiding this comment.
The PR description has become outdated following the latest threading changes. The original race diagram remains useful for illustrating the pre-fix behavior, but the Modifications and Verification sections no longer describe the current implementation. Changes include: flow updates now run on the dispatchMessagesThread, executor rejection intentionally leaves credits pending, and queued updates use exact Consumer-instance membership to avoid applying stale credit to a replacement Consumer with equal attributes. A new SharedDispatcherFlowThreadingTest has also been added, along with JFR and mixed-stress validation.
Since this PR changes the threading model and is intended for backports, please refresh the description to reflect the final design before merging.
There's a parse error in rendering this diagram. |
lhotari
left a comment
There was a problem hiding this comment.
The three new commits move Flow processing off the connection EventLoop instead of inlining it there, which is a different design from the one discussed earlier on this PR. The reasoning and the JFR numbers for that reversal are convincing on their own terms, but it is a broker-wide threading change riding on a permit-accounting fix, so it deserves an explicit decision rather than arriving as a follow-up commit. I did not find a correctness defect in the accounting itself.
What I checked and found sound at f81fcae0:
- The permit invariant is still exact.
Consumer.flowPermitsraisesmessagePermitsandpendingDispatcherFlowPermitstogether underflowPermitAccountingLock(service/Consumer.java:974-982), so the removal balance is unchanged until the dispatcher task runs; the task settles the pending value and creditstotalAvailablePermitsunder the dispatcher monitor, for the exact connected instance (PersistentDispatcherMultipleConsumers.java:324-333); removal subtracts the resulting balance (PersistentDispatcherMultipleConsumers.java:262-265). No path holdsflowPermitAccountingLockwhile acquiring the dispatcher or subscription monitor, so the ordering is one-way. - Leaving the permits pending on
RejectedExecutionExceptionis the right call. Thebroker-topic-workersexecutor is built with no queue limit (BrokerService.java:389-391), so it rejects only once it has been shut down during broker shutdown, not under load. - The identity membership check is load-bearing, not merely defensive.
Consumer.equalsisconsumerIdpluscnx.clientAddress()(service/Consumer.java:1184-1193), and in the replacement race the set holds the new instance, socontains()would have let a stale Flow task credit permits for a consumer that had already been removed.testQueuedFlowDoesNotApplyToEqualReplacementConsumerpins exactly that. - No new recursion into
readMoreEntries(). The Flow call is a top-level dispatch-lane task, the modern dispatcher already callsreadMoreEntries()synchronously from its send/read loop on the same lane, and the recursion-sensitive branches still usereadMoreEntriesAsync(). - All five earlier threads still hold at this head, including the two whose test mechanism changed:
testRemainingConsumerCanContinueAfterFlowAndCloseRacestill follows the subscription-then-dispatcher order and still waits for the connection-map removal, and the queued-Flow ordering is now forced with a dispatch-lane barrier rather than by holding the dispatcher monitor, which is strictly safer now that the monitor no longer gates the Flow task.
The consequences of the new placement, and the shape I would like the read triggering to take, are in the inline comments below.
Scope and backport. The labels put this on 4.0.14 and 4.2.5. The accounting fix is contained -- Consumer bookkeeping plus the two removeConsumer call sites -- while the new commits change which executor every persistent Shared and Key_Shared consumerFlow runs on, and how reads are triggered for all of them. Would you consider splitting the threading change into its own PR, so what lands on the maintenance branches is only the accounting fix?
|
Thanks for the JFR numbers -- they do make the case against inlining the accounting on the EventLoop, and I agree with that much. Two things I would still like to settle before this goes in. Where I am not with you yet is the replacement. The wait is relocated rather than removed: And the read trigger on the Flow path no longer goes through Separately, on scope: the release labels put this on 4.0.14 and 4.2.5. The accounting fix is contained, the threading change is not. I would rather see them split so the maintenance branches take only the accounting fix -- would that work for you? |
Restore the existing executor and read-trigger behavior while retaining exact Consumer identity and rejected-submission accounting coverage. Remove the threading-specific tests and keep the behavior matrix for modern/classic Shared and Key_Shared dispatchers. Assisted-by: OpenAI Codex
|
Thanks for the detailed review and for calling out the scope, backport, and read-trigger concerns. I updated the PR in I removed the Flow migration to The PR now retains only the accounting-related changes and behavior coverage, including scoped pending accounting, signed balances, exact Consumer-instance matching, and rejected-executor handling. I also updated the PR description and corrected the Mermaid diagram to match the final implementation. Thanks again for the careful review and the concrete suggestions. |
lhotari
left a comment
There was a problem hiding this comment.
The accounting-only revision addresses my remaining scope and read-trigger concerns. The signed pending-permit accounting and regression coverage are retained. I found no further blocking issue in d72e6c7a9c8.
|
Thanks for narrowing this to the accounting fix. The restored executor placement and modern async read trigger match the scope I requested. The broader read-conflation work can be considered separately. |
The existing churn suite asserts that dispatch CONTINUES across a mid-drain consumer close. That is green on a broker whose permit ledger is leaking: `readMoreEntries` has fallen back to `max(totalAvailablePermits, firstAvailableConsumerPermits)` since 2021, so a negative aggregate is masked entirely. "Messages kept flowing" therefore proves nothing about issue #414's root cause. This suite asserts on the numbers the broker reports about itself instead: the dispatcher's `totalAvailablePermits` never negative, every per-consumer `availablePermits` inside `[0, receiver_queue_size]`, no negative `unackedMessages`, and a backlog that still reaches zero. The aggregate is exposed by no admin endpoint — measured absent from `…/stats`, `/admin/v2/broker-stats/topics`, `…/internalStats` and `/metrics` on 4.2.4 — so the suite raises the `PersistentDispatcherMultipleConsumers` debug logger inside the container and parses the counter out of the broker log. Parsing zero observations fails as an infrastructure fault rather than passing. It also fixes a real trap for 5.x: `WaitFor::message_on_stdout("Created namespace public/default")` times out there because 5.x logs it to stderr. It polls the namespace list instead, like `e2e_scalable_topic.rs`. `#[ignore]`d, which ADR-0046 otherwise forbids, and this is the documented exception. Every other e2e test asserts something about magnetar; this one asserts that apache/pulsar#26416 is absent, which no client behaviour can change and which no generally-available image has fixed — Docker Hub stops at 4.0.13 and 4.2.4, and the fixed 4.0.14 / 4.2.5 are unpublished. Left running on the `latest` default it would be permanently red in CI on a defect this repository cannot repair, which is the failure mode that turns a check into noise. The attribute carries the reason and the reproduce command, and the `latest` default is kept so it goes green by itself the day a fixed image ships. Verified across repeated runs rather than one cell, because the result is not a clean red/green: 4.0.4 red 3/3 (min -10 to -14), 4.2.4 (== `latest`) red 8/8 (min -4 to -18), and 5.0.0-M2 -- the only image carrying apache/pulsar#26289 + apache/pulsar#26422 -- green in 11 of 12, with the twelfth reaching `min=Some(-2) negatives=1`. The merged upstream fixes therefore shrink the leak by about two orders of magnitude without closing it, which is why the assertion is written against the invariant and not against a known-good image. The full table is in `docs/testing.md`. Signed-off-by: Florentin Dubois <[email protected]>
…loss during frequent reconnects on Shared subscriptions (apache#26289) (cherry picked from commit 12b6497)
…loss during frequent reconnects on Shared subscriptions (apache#26289) (cherry picked from commit 12b6497)
…loss during frequent reconnects on Shared subscriptions (apache#26289)
Fixes #26288
Motivation
For persistent Shared and Key_Shared subscriptions, a Flow command updates the Consumer permit counter before the corresponding update is applied to the dispatcher total. Previously, consumer removal could run while that dispatcher update was still queued.
The following diagram shows the pre-fix failure. Time flows downward.
sequenceDiagram participant IO as Connection EventLoop participant C as Removed Consumer participant W as Queued Flow Task participant D as Shared Dispatcher IO->>C: Receive 1000 additional permits C->>C: Available permits become 1010 C-->>W: Queue dispatcher update IO->>D: Remove consumer before queued update runs D->>D: Dispatcher total becomes negative W->>D: Process queued update D->>D: Ignore update from disconnected consumerFor example, if the removed consumer had 10 permits and the dispatcher total was 20, removing the consumer after its local permit count increased by 1000 would subtract 1010. The queued update would then be ignored because the consumer was already gone, leaving the dispatcher total at -990 instead of 10.
This negative permit drift does not necessarily stop consumption. It reduces the effective permits used to size reads and can cause persistent throughput degradation after frequent reconnects.
Modifications
messagePermits - pendingDispatcherFlowPermits.Accounting model and scope
When a Flow command is accepted, the Consumer updates
messagePermitsandpendingDispatcherFlowPermitstogether under a per-Consumer accounting lock. The existing asynchronous dispatcher task settles the pending value before applying or ignoring the update. Consumer removal subtracts onlymessagePermits - pendingDispatcherFlowPermits, so permits not yet represented in the dispatcher total are excluded.Under the dispatcher monitor, the maintained invariant is:
The balance is intentionally signed and preserves Java
intwrap behavior. It may legitimately be negative when dispatch consumes permits before a queued Flow update is applied.This PR does not change executor placement or read-trigger behavior. Flow processing remains on
BrokerService.executor()in both dispatchers. The modern dispatcher retains the deduplicatedreadMoreEntriesAsync()path, and the classic dispatcher retains its existing direct-read path. No executor or thread pool is added.Verifying this change
Targeted validation completed on
home3against final commitd72e6c7a9c8:SharedDispatcherPermitAccountingTest: 25/25 passed../gradlew :pulsar-broker:checkstyleMain :pulsar-broker:checkstyleTest: passed../gradlew quickCheck: passed.The tests can be run with:
Does this pull request potentially affect one of the following parts:
The change adds a short per-Consumer accounting critical section for persistent Shared and Key_Shared subscriptions. No callback, dispatcher monitor, or I/O operation is executed while holding that accounting lock. No user-facing client/admin API, wire protocol, persisted schema, configuration/default, or metric schema is changed; the added Consumer helpers are broker-internal.