Repository navigation
[fix][ml] Fix active cursor retention and deferred-update synchronization - #26511
Merged
Merged
Conversation
…tion Release removed cursor references immediately and compact stale nodes in batches without flushing lazy position updates. Track queue membership to prevent duplicate entries during reactivation. Fix deferred-state mutation under a shared read lock. Use optimistic reads for clean slowest-position queries and exclusive locking when flushing. Add regression tests and cursor churn/read-only JMH workloads. Assisted-by: OpenAI Codex
lhotari
requested review from
Technoboy-,
dao-jun,
david-streamlio,
merlimat and
nodece
September 9, 2026 08:59
10 tasks
Remove the tail-only early return from incremental position updates so joining and leaving a shared position group updates cursor ranks. Add regression coverage for a moving cursor that is removed and reactivated, tail group transitions, and subsequent node compaction. Assisted-by: OpenAI Codex
Make add() clear pending removal state when an existing cursor receives a real position, matching updateCursor(). Add regression coverage for both rebuild and incremental flushing and both ways of untracking. Assisted-by: OpenAI Codex
Keep the tail counter when its old and new positions are separate from the preceding group. Preserve counter updates when joining or leaving a shared position group, and add a focused tail-seek benchmark. Assisted-by: Codex
nodece
approved these changes
Sep 9, 2026
lhotari
marked this pull request as draft
September 9, 2026 11:10
Add regression tests for exclusive pending-position flushes and appending after compaction removes the tail. Strengthen the test-only state checker to validate cursor references, backward links, the tail, and tracked node count without changing production operation paths. Assisted-by: Codex
dao-jun
approved these changes
Sep 9, 2026
lhotari
marked this pull request as draft
September 9, 2026 12:12
Accumulate pending nodes in an ObjectArrayList and preserve descending old-position processing during incremental flushes. Use the small-array sort for short batches, skip sorting ordered larger batches, and use in-place quicksort otherwise to avoid temporary merge buffers. Full rebuilds skip pending-list sorting. Cover batch boundaries, repeated updates, new cursors and shared position groups. Add a grouped-advance benchmark. All 59 cursor cases and quickCheck pass; final paired JMH comparisons completed without a measured throughput regression. Assisted-by: Codex
lhotari
marked this pull request as ready for review
September 9, 2026 13:26
1 of 11 tasks
dao-jun
pushed a commit
to ascentstream/pulsar
that referenced
this pull request
Sep 20, 2026
…tion (apache#26511) (cherry picked from commit 72632e5)
Radiancebobo
pushed a commit
to Radiancebobo/pulsar
that referenced
this pull request
Oct 8, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Replacement for #26506.
Motivation
ActiveManagedCursorContainerImpldefers cursor ordering until a position query needs it. WithcacheEvictionByExpectedReadCountenabled, cursor creation and deletion can continue without such queries. Removed nodes then accumulate inpendingPositionUpdatesor remain in the ordered list andpendingRemovedCursors, retaining their cursor objects and associated object graphs.This PR replaces the approach in #26506 with immediate cursor-reference release and bounded, batched node compaction. The ten-event flush in #26506 counts cursors removed before their first position flush; already-tracked removals are outside that cleanup trigger. Preserving lazy position updates is important to the efficiency of this data structure.
This also fixes a pre-existing synchronization bug:
getSlowestCursorPosition()held a read lock while callingprocessPendingPositions(), which mutates the pending-update collection, linked list, and counters. Concurrent readers could therefore flush the same state simultaneously. The internal state-checking method had the same locking problem. This is API-correctness hardening: under the current broker configuration paths, this implementation is selected whencacheEvictionByExpectedReadCountis enabled, and the broker skips its slowest-position lookup in that mode. The production rank-query path already used the write lock.Modifications
max(64, live cursor count)removals, or immediately when the container becomes empty. Prune the pending-update list and unlink removed list nodes in linear time, repairing shared counters in one pass. Compaction preserves pending position updates and does not sort or flush them.pendingPosition, so cancellation and reactivation cannot add the same node repeatedly.ObjectArrayListand sort only before incremental flushes. Process descending old positions to reduce list traversal when cursors advance together. Small batches use the standard sort; larger batches skip sorting when already ordered and otherwise use in-place quicksort to avoid temporary merge buffers. Full rebuilds skip pending-list sorting. Fastutil is already a dependency.add()cancel pending untracking when an existing cursor receives a real position, matchingupdateCursor(). This is a defensive consistency fix: production activation currently guards against adding an already-active cursor.ActiveManagedCursorContainerBenchmarkwith cursor-churn, read-only slowest-position, and tail-seek workloads. AddActiveManagedCursorContainerCohortBenchmarkfor grouped cursor advances.Verifying this change
Local validation passed:
The two test classes now run 59 cases. The 11 retention regression cases fail on the unpatched implementation with only read-only testing accessors added. Coverage includes removal before and after position tracking, bounded retention without queries, preserving deferred updates and shared counters during compaction, repeated reactivation, replacement cursor identity, iterator behavior, and cached reads completing while another reader holds the lock. Two additional regression cases fail before the tail-counter fix and cover reactivation after a pending move, joining and leaving position groups, and subsequent compaction. Four further regression cases cover cancelling pending untracking through
add()on both rebuild and incremental paths; all four fail before the defensive fix. Two further tests protect the pending-flush write lock and compaction tail repair. Both were mutation-tested: replacing the flush write lock with a read lock fails the first test; removingtail = previousfrom compaction fails the second. The additional structural checks run only in the test helper, not on production paths. A further 24 parameterized cases verify batches around the sorting threshold, ascending/reverse arrival order, mixed new and tracked nodes, coalesced repeated updates, shared positions, and a second flush to check pending-state reuse.Performance measurements
The final hybrid implementation at
108223dba7b4was measured against this PR's previous PriorityQueue implementation atc80d72a4db89, plus a fresh churn comparison against #26506. All 18 configurations completed and passed their benchmark state checks. Both sides use the same final benchmark jar and dependencies; only implementation classes are overridden for the comparison variants.Environment: macOS arm64, Corretto 21.0.11, DEFAULT implementation, 256 MiB heap, GC profiler. Each configuration uses two forks, two 1-second warmup iterations and three 1-second measurement iterations per fork. Variants run sequentially, alternating order between parameter combinations. Tables show millions of updates/operations per second ± JMH-reported 99.9% error; percentages compare means. These are short local measurements, not production performance guarantees. Linux x86_64 validation is still needed.
Final hybrid versus the previous PriorityQueue implementation
Query ratio is updates per rank query; zero disables rank queries.
randomSeekingForward01uses one thread andrandomSeekingForward10uses ten.cohortis the newActiveManagedCursorContainerCohortBenchmark.advanceCohort: it advances the fastest 32 cursors together in ascending old-position order, then queries rank. Both population sizes force incremental processing, and@OperationsPerInvocation(32)normalizes throughput and allocation per cursor update. This deliberately ordering-sensitive workload demonstrates the traversal benefit; it is not a measured broker traffic distribution.randomSeekingForward01randomSeekingForward01randomSeekingForward01randomSeekingForward01randomSeekingForward10cohortcohortcursorChurnCursor churn versus #26506
At 500 active cursors with no rank queries,
cursorChurnmeasures 15.09 ± 0.56 M ops/s for #26506 versus 16.66 ± 0.48 M ops/s for this PR, a +10.4% change in mean throughput. #26506's implementation is pinned tobbf855c72c33and overrides only implementation classes in the same benchmark jar.Why this sorting strategy
FIFO experiments improved several random-seek workloads but lost roughly half their throughput when cursors advanced together. Sorting only during incremental flushes preserved descending-old-position processing, but plain ArrayList.sort introduced temporary-buffer allocations for larger batches. The final hybrid uses the standard small-array sort below 32 entries on JDK 21; for larger batches it skips sorting if already ordered, otherwise uses Fastutil's in-place quicksort. Full rebuilds skip pending-list sorting. Separate paired experiments against plain deferred sorting reduced allocation from 58.37 to 46.37 B/update for large random batches and from 80.00 to 47.50 B/update for the 500-cursor cohort. Fastutil is already a dependency.
The read-only slowest-position benchmarks remain API coverage; this implementation's slowest-position lookup is not on the current broker path.
Reproducing the current implementation's workloads
Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes