Repository navigation
[improve][broker] Avoid redundant read-completion executor handoffs for Exclusive/Failover - #26619
Conversation
Allow self-dispatching read callbacks to complete on the current thread while retaining the Exclusive/Failover dispatcher executor boundary. Add callback affinity and failure coverage, a handoff benchmark, and reproducible YAML profiling scenarios. Assisted-by: Codex
void-ptr974
left a comment
There was a problem hiding this comment.
Nice improvement. Removing the unnecessary managed-ledger executor handoff brings a significant performance gain for Exclusive/Failover consumers, especially under heavy producer load.
I also like that this is implemented as an opt-in mechanism rather than changing the callback threading behavior globally. This keeps the change well scoped and also provides a useful foundation for similar optimizations in other consumer paths.
LGTM.
@void-ptr974 The next change #26620 makes it used for all cases since it's needed for all use cases. |
Motivation
With many producers publishing to one persistent topic, publish completions can saturate the topic's ordered managed-ledger executor. Consumers then cannot keep up with the publish rate, causing a significant subscription backlog. Adding broker cores does not eliminate this per-topic serialization bottleneck; faster cores may raise the saturation point, but leave consumer dispatch dependent on the same saturated executor.
For Exclusive and Failover subscriptions, a cached read initiated by the dispatcher sends its successful completion to the managed-ledger executor before
whenCompleteAsyncsends the continuation back to the dispatcher executor. Even confirmed messages already in the entry cache must wait behind publishing work. This PR removes the first handoff for this self-dispatching callback, so consumers can keep up without scaling up the broker to compensate for the avoidable queue delay.Modifications
ReadEntriesCallback.canExecuteOnAnyThread()capability, defaulting tofalse, and a corresponding future-adapter overload.whenCompleteAsynccontinuation. Failure and compacted-read paths retain their existing behavior.tests/performance/scenarios/, with a public reproduction guide.Correctness and ordering
Confirmed positions remain the read-visibility boundary.
ManagedLedgerImpl.internalReadFromLedgercaps active-ledger reads at the managed ledger'slastConfirmedEntry, which may lag the BookKeeper handle's LAC (last add confirmed); closed-ledger reads use BookKeeper's LAC. The topic'smaxReadPositionis an additional bound, preserving transaction visibility.OpAddEntryinserts an entry into the cache and advances the volatilelastConfirmedEntryonly after a successful BookKeeper add. The managed-ledger limit ensures its successful-add bookkeeping is visible before entries become readable. Completing a read on another thread does not expose an entry before that confirmation. The producer receiving its acknowledgement is a separate network event and need not precede consumer delivery.Ordering follows the cursor's ordered reads and the dispatcher's read lifecycle, rather than the identity of the successful-completion thread:
OpReadEntryassembles entries in cursor order and advances the read position before invoking the final successful callback. The change affects that callback's scheduling, not which positions are read or their order. Existing partial-batch completion behavior is preserved.CompletableFuture.whenCompleteAsync(..., executor)still schedules processing on the dispatcher executor, preserving the asynchronous boundary and state ownership. Inline future completion does not recursively invoke the dispatch loop on the completing thread. The volatile confirmed position and future/executor publication retain cross-thread visibility.Verifying this change
Local validation passed: full
./gradlew spotlessCheck checkstyleMain checkstyleTest; 28 managed-ledger/cache cases including the upstream cached-read recursion bound and real-BookKeeper ordering/thread-affinity tests; 5 broker dispatcher cases; and 18 YAML configuration cases. The tests use the sharedManagedLedgerTestUtil.rawEntryConfig()helper.A fresh baseline/candidate integration comparison after rebasing onto upstream master used 500 producers on 500 isolated v4 client connections, one Exclusive consumer, one persistent topic, 12 million unbatched 128-byte messages, unrestricted publishing and 40 outstanding sends per producer. Both runs consumed all 12 million messages with zero failed ACKs, one broker startup and no OOM. Only this scenario was measured for this comparison; the other YAML variations are available for follow-up runs.
Steady dispatch uses topic
msgOutCounterdeltas while all 500 producers are active. The window starts after the first 20 seconds with all producers connected and ends before producers start disconnecting; only complete intervals with 500 publishers at both endpoints count. The selected windows were 86 seconds for baseline and 67 seconds for candidate. Steady ingress usesmsgInCounterover the same intervals. These measure broker delivery/publishing, not client acknowledgement.Whole-run throughput uses each
pulsar-perfclient's own aggregate timer and includes ramp-up and backlog drain. The baseline can therefore show a much higher whole-run consumer rate than its delivery rate during publishing. Sampled maximum backlog is the maximum across all collected topic snapshots, with a configured five-second polling delay (actual intervals are about six seconds including collection work); it can miss peaks between samples. Lower backlog and CPU are improvements; throughput percentages are increases.Whole-run consumer throughput improved 85.2%, and steady dispatch kept up with ingress instead of accumulating millions of messages.
CPU, lock and allocation profiles were collected with
event=cpu,interval=10ms,lock=0,alloc=2m,jfrsync=profile, without wall-clock sampling, and analyzed with Jafar MCP diagnosis and stack profiles. The handoff-only JMH benchmark (two forks, 3 x 1s warmup and 5 x 1s measurement, GC profiler) measured 9.53 ± 0.44 → 4.47 ± 0.16 microseconds/op and 312 → 232 bytes/op. It models executor/future overhead, not storage IO or end-to-end capacity.This is one sequential run pair on an 8-core/16-thread i9-9980HK with performance power settings and SMT enabled. Both runs experienced thermal throttling (steady maxima 98°C/95°C), so exact capacity differences need repetition under controlled cooling. Broker allocation per second is not an equal-work comparison because the candidate delivers substantially more messages during publishing.
Related BookKeeper batch-read queueing
Further investigation identified a separate redundant completion queue in BookKeeper, reported as apache/bookkeeper#4893.
LedgerHandle.batchReadUnconfirmedAsync()queues its future continuation even when the batch-read response already completes on the ledger's worker. Fixing that queue alone does not remove the cached-read completion handoff addressed by this PR.In a separate controlled follow-up with 500 producer connections and one Exclusive consumer, removing only that BookKeeper queue increased steady dispatch from 5,435 to 9,894 messages/s, while ingress remained 105,314 messages/s with the change. The mean BookKeeper read-future duration fell from 169.4 to 92.7 ms, but the first response-to-worker queue still took 84.7 ms. The consumer still accumulated a sampled peak backlog of 10,432,707 messages. Both runs consumed all 12 million messages with zero failed ACKs, no broker restart and no OOM.
Thus the measured BookKeeper-only prototype reached about 10k messages/s during publishing, still far below the publish rate. This is a separate comparison using a dispatcher variant with coalesced read requests, channel-writability-based resumption, and 256/128 KiB channel watermarks; it is not another measurement of the baseline/candidate pair above. Both runs reached 98°C and throttled, so treat the precise rates as scenario-specific. This PR lets already-confirmed cache hits complete without the extra managed-ledger queue hop.
Does this pull request potentially affect one of the following parts: