Skip to content

[improve][broker] Optimize consumer selection for shared subscriptions - #26593

Merged
merlimat merged 3 commits into
apache:masterfrom
lhotari:lh-improve-shared-consumer-select
Sep 15, 2026
Merged

merlimat merged 3 commits into
apache:masterfrom
lhotari:lh-improve-shared-consumer-select

Conversation

@lhotari

@lhotari lhotari commented Sep 15, 2026

Copy link
Copy Markdown
Member

Motivation

Shared subscriptions perform priority checks during consumer selection even when every consumer has the same priority. With mixed priorities, round-robin wraparound can also rescan higher-priority consumers to locate the current priority group.

Reducing this work lowers selection overhead while preserving priority ordering and round-robin behavior.

Modifications

  • Extract selection logic into ConsumerPrioritySelector, allowing isolated unit tests and JMH benchmarks.
  • Track consumer counts by priority using Int2ObjectMap<MutableInt>. Remove entries when their counts reach zero and skip priority checks when only one priority remains.
  • Reuse priority-group boundaries and eliminate redundant list reads during mixed-priority selection.
  • Synchronize selection and membership updates on the dispatcher monitor. Preserve cleanup of inconsistent consumer list/set state by rebuilding counts during bulk-removal recovery.
  • Keep the refactoring, benchmark, and optimization in separate commits for reproducible comparisons.

Verifying this change

  • Make sure that the change passes the CI checks.

Added tests cover:

  • Selection and availability-check order across 262,144 priority-group, availability-mask, and cursor combinations.
  • Three complete round-robin cycles with 50 consumers, followed by exhaustion and recovery, at zero and nonzero priority.
  • Membership changes, priority-count transitions, equal consumer identities, and bulk-removal recovery.

Local validation passed:

  • ConsumerPrioritySelectorTest
  • The four existing persistent dispatcher priority tests, with quarantine exclusions cleared and retries disabled
  • Selected dispatcher add/remove and cleanup tests
  • ./gradlew quickCheck
  • ./gradlew :microbench:shadowJar

JMH comparison

Compared baseline 2d540fa0853 with optimization 9c88617c944 using identical benchmark sources and settings across 26 parameter combinations.

Environment: Apple M3 Max, macOS arm64, Corretto 25.0.4, JMH 1.37; two forks, three 500-ms warmup iterations, five 500-ms measurement iterations, one thread, and the GC profiler. Builds and tests did not overlap measurements.

Values below are mean nanoseconds per selection. Ranges cover priority offsets zero and three; speedups compare matching cases.

Consumers Priority levels Availability Baseline ns/op Optimized ns/op Speedup
50 1 All ready 3.13–3.96 2.82–2.84 1.10–1.40×
50 1 Sparse 17.36–21.28 8.52–8.94 2.04–2.38×
50 1 All blocked 90.27–91.28 46.20–49.91 1.81–1.98×
16 4 All blocked 43.08–44.12 20.89–21.03 2.05–2.11×
256 4 All blocked 609.33–614.69 292.88–296.06 2.06–2.10×

No GC occurred, and measured allocation was approximately zero. The 256-consumer mixed-priority FLOW case at offset zero had a wide confidence interval and remains inconclusive.

These measurements isolate selection cost and exclude dispatcher locking, membership updates, and message processing. Linux x86_64 confirmation is needed before extrapolating to broker throughput.

Does this pull request potentially affect one of the following parts:

  • Dependencies (add or upgrade a dependency)
  • The public API
  • The schema
  • The default values of configurations
  • The threading model
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

Selection and membership helpers explicitly synchronize on the existing dispatcher monitor so the consumer list and priority counts remain consistent. Existing dispatch call sites already hold that monitor.

Preserve selection and availability-check order while isolating list traversal for benchmarking. Keep the dispatcher cursor and lifecycle checks in place.

Assisted-by: Codex
Measure ready, sparse, lower-priority, blocked and changing-permit consumers with persistent round-robin state. This revision is the unoptimized baseline.

Assisted-by: Codex
Track MutableInt priority counts with synchronized membership updates, removing entries when their count reaches zero. Use availability-only round-robin traversal for one priority and reuse group boundaries for mixed priorities. Preserve the bulk-removal repair path for inconsistent dispatcher membership.

Validate with JMH comparisons against the benchmark-only parent, including 50 consumers at one priority, selection-order and membership tests, dispatcher recovery tests, explicitly enabled priority tests, and quickCheck.

Assisted-by: Codex
@merlimat
merlimat merged commit e7cbb33 into apache:master Sep 15, 2026
44 checks passed
dao-jun pushed a commit to ascentstream/pulsar that referenced this pull request Sep 20, 2026
dao-jun pushed a commit to ascentstream/pulsar that referenced this pull request Sep 20, 2026
dao-jun pushed a commit to ascentstream/pulsar that referenced this pull request Sep 20, 2026
@lhotari lhotari added this to the 5.0.0 milestone Oct 1, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants