Skip to content

[improve][broker] PIP-379: Remove the classic Shared and Key_Shared dispatcher implementations - #26687

Merged
lhotari merged 5 commits into
apache:masterfrom
lhotari:lh-remove-legacy-dispatchers
Sep 22, 2026
Merged

lhotari merged 5 commits into
apache:masterfrom
lhotari:lh-remove-legacy-dispatchers

Conversation

@lhotari

@lhotari lhotari commented Sep 22, 2026

Copy link
Copy Markdown
Member

PIP: 379 (pip/pip-379.md)

Motivation

PIP-379 replaced the Key_Shared "recently joined consumers" machinery with the draining-hashes
tracker, and the matching work replaced the Shared dispatcher. Both landed in 4.0.0 LTS
(Oct 2024) and became the default in that release. The pre-4.0 implementations were kept alive
behind two feature flags purely as a rollback escape hatch:

  • subscriptionKeySharedUseClassicPersistentImplementation
  • subscriptionSharedUseClassicPersistentImplementation

5.0 is the right release to delete them.

The classic path was already an escape hatch, not a supported alternative

It shipped in 4.0.0 with every signal of a temporary fallback:

  • Both flags default to false, so no default deployment has used the classic path since 4.0.0.
  • Neither flag appears in conf/broker.conf or conf/standalone.conf — they were never advertised
    as a supported knob.
  • The classic-only stats fields ConsumerStats.getReadPositionWhenJoining() and
    getKeyHashRanges() were marked @Deprecated in [improve][broker][PIP-379] Add observability stats for "draining hashes" #23429 back in Oct 2024, i.e. they shipped
    deprecated in 4.0.0 and have carried that warning through 4.0, 4.1 and 4.2.
  • pip/pip-379.md already lists consumersAfterMarkDeletePosition as removed.

Users have therefore had two full LTS lines of notice, with the modern implementation as the
default the whole time.

PIP-379 has been the default for two years; the classic path still has the original defects

The draining-hashes implementation has been the default for every deployment on 4.0+ for roughly
two years, across three release lines. The classic implementation still has the problems PIP-379
documents in its Motivation section:

  1. an ordering contract stated in terms of its own solution rather than as a guarantee;
  2. incomplete fulfilment of that contract — a key can be outstanding in more than one consumer
    during consumer changes (#23307);
  3. hard-to-diagnose blocking;
  4. unnecessary blocking: the classic path blocks delivery for all messages when any hash range
    is blocked, even when other keys could be dispatched independently;
  5. no visibility into why dispatching is stuck — the modern path exposes drainingHashes,
    drainingHashesCount, drainingHashesClearedTotal and drainingHashesUnackedMessages;
  6. the "recently joined consumers" mechanism itself, which PIP-379 calls out as overly complex.

It also carries #21199 (Key_Shared subscription
gets stuck after consumer reconnects). Keeping the flag keeps a supported way to opt back into
these defects.

The classic path is excluded from the PIP-430 cache improvements

PIP-430's expected-read-count eviction (cacheEvictionByExpectedReadCount, default true)
relies on the dispatcher incrementing an entry's expected read count when that entry is pushed onto
the replay queue, so the entry is de-prioritised for size-based eviction and is more likely to still
be cached when the Key_Shared consumer can finally take it (PIP-430, "Detailed Design").

That call — entry.getReadCountHandler().incrementExpectedReadCount() — exists only in
PersistentStickyKeyDispatcherMultipleConsumers. The classic dispatcher never had it. So on 5.0 a
user who sets the classic flag would silently get materially worse cache hit rates on exactly the
replay traffic the cache improvement was designed for, with nothing pointing at the flag as the
cause. Features built on top of PIP-379 will keep widening this gap.

Maintenance cost, concretely

The two classes are ~2000 lines of duplicated dispatcher logic, plus a classic/modern branch
threaded through PersistentSubscription, MessageRedeliveryController and the admin stats
surface, and a @Factory that started a second broker for the largest Key_Shared test class.

The clearest evidence is the commit this branch is based on:
#26677 "Bound classic Key_Shared dispatcher replay
queue look-ahead". It retrofits a look-ahead bound onto the classic dispatcher by reaching into the
modern class for PersistentStickyKeyDispatcherMultipleConsumers.getEffectiveLookAheadLimit(...),
and even then only approximates it — its own comment concedes "a bounded look-ahead limit with an
acceptable one-batch overflow". That is the ongoing tax: a capability designed for the modern
dispatcher has to be re-derived, approximately, for a path nobody runs by default, and every future
fix in this area faces the same choice between doing the work twice or leaving the two paths
divergent.

Deleting it removes that tax for the whole 5.0 line.

Modifications

Five commits, ordered so the tree compiles and tests pass at each one. Deletion cannot come first:
eight test files call the config setters, six name the classic classes, and ConsumerStatsTest
asserts an exact JSON field set — a runtime dependency javac would not catch. The tests are
therefore de-parameterized first.

  1. [improve][test] Stop exercising the classic Key_Shared implementation in tests — delete
    org.apache.pulsar.tests.KeySharedImplementationType and unwind its four users (the @Factory,
    the extra constructors, the inject-only currentImplementationType provider, and the
    prependImplementationTypeToData wrapping). testReadAheadLimit and
    testOrderingAfterReconnects lose their skipIfClassic() and now always run.
  2. [improve][test] Stop parameterizing Shared-dispatcher tests on the classic implementation —
    ConsumerStatsTest, SharedDispatcherPermitAccountingTest,
    SharedSubscriptionUnackedMessagesAccountingTest and
    PersistentDispatcherMultipleConsumersTest.
  3. [improve][broker] Remove the classic Shared and Key_Shared dispatcher implementations —
    delete PersistentDispatcherMultipleConsumersClassic,
    PersistentStickyKeyDispatcherMultipleConsumersClassic, their dedicated tests and
    NonEntryCacheKeySharedSubscriptionV30Test; delete both ServiceConfiguration flags; collapse
    the branches in PersistentSubscription.
  4. [improve][broker] Drop the dispatcher seams that only the classic implementations used —
    isClassic(), StickyKeyDispatcher.getRecentlyJoinedConsumers(), and the
    MessageRedeliveryController isClassicDispatcher field, 2-arg constructor and
    containsStickyKeyHashes(Set<Integer>) overload.
  5. [improve][admin] Remove the classic-only consumer and subscription stats fields — the only
    commit with a public API break, kept as one reviewable diff. See Upgrade impact.

Net: 28 files, +208 / −4031 (production code: +13 / −2173).

The PIP-486 ksm.isEntryBucketDispatch() arm in PersistentSubscription is kept, as is
AbstractPersistentDispatcherMultipleConsumers — it is the parameter type of the user-pluggable
DelayedDeliveryTrackerFactory.newTracker(...) SPI, the narrowing type BrokerService uses for
blockedDispatchers, and the type seven test files mock.

Deliberately out of scope: the PIP-322 classic dispatch rate limiter
(DispatchRateLimiterClassicImpl, DispatchRateLimiterFactoryClassic,
dispatchRateLimiterFactoryClassName). Different feature, first shipped in 4.1.0 — a much shorter
deprecation life. Separate PR.

Upgrade impact

Stale config is ignored, not rejected. A leftover
subscriptionKeySharedUseClassicPersistentImplementation=true in broker.conf does not fail
startup — FieldParser.update iterates declared fields, never the properties — and surfaces only in
the configOverrides attribute of the "Messaging service is ready" startup log. Both flags were
dynamic = true, so a stale value may also sit in /admin/configuration; BrokerService already
logs and skips unknown dynamic keys. The broker simply uses the modern dispatcher, which is what it
already did for every default deployment since 4.0.0.

One JSON-visible stats change. SubscriptionStats.consumersAfterMarkDeletePosition was
initialized to an empty LinkedHashMap, so persistent subscriptions emit
"consumersAfterMarkDeletePosition": {} today and will stop emitting it. The other two removed
fields are null on the modern path and the admin mapper uses Include.NON_NULL, so their removal
is JSON-invisible. None of the three is exposed via Prometheus metrics. This is a deliberate 5.0
stats-shape change.

Why consumersAfterMarkDeletePosition is removed rather than renamed or kept as an empty map.
It was never a general "consumers past the mark-delete position" statistic that some other mechanism
could repopulate — it was the classic dispatcher's recently-joined-consumers map under a name that
described its eviction rule:

  • Its own javadoc contradicted its name. Both SubscriptionStats and SubscriptionStatsImpl
    carried the line "This is for Key_Shared subscription to get the recentJoinedConsumers in the
    Key_Shared subscription." The documented subject is recentJoinedConsumers.
  • The name described a retention rule, not the data. The key is a consumer name; the value is that
    consumer's cursor read position at the moment it joined — the same value surfaced as the
    consumer-level readPositionWhenJoining. "After mark delete position" referred only to the
    pruning invariant in removeConsumersFromRecentJoinedConsumers(), which dropped an entry once the
    mark-delete position passed the join position. The name conveys neither the key nor the value.
  • It had no meaning outside the feature. Its sole writer was PersistentSubscription, fed
    exclusively by StickyKeyDispatcher.getRecentlyJoinedConsumers(), whose only overrider was the
    classic dispatcher. PIP-282 renamed the consumer-level sibling (readPositionWhenJoining →
    lastSentPositionWhenJoining) but left this one alone, and PIP-379 then declared it removed.

The equivalent modern observability is richer: drainingHashes, drainingHashesCount,
drainingHashesClearedTotal, drainingHashesUnackedMessages and keyHashRangeArrays.

Verifying this change

This change is already covered by existing tests.

The end-to-end proof that the removal is behaviourally inert on the default path is
ConsumerStatsTest.testConsumerStatsOutput, which asserts the exact field set of the
/admin/v2/.../stats JSON: after commit 2 it pins the modern Key_Shared stats shape, and after
commit 5 it proves nothing else shifted.

Run locally, all green:

  • ./gradlew quickCheck sanityCheck verifyTestGroups
  • :pulsar-broker:test --tests "org.apache.pulsar.client.api.KeySharedSubscription*" --tests "org.apache.pulsar.client.impl.KeySharedSubscriptionMaxUnackedMessagesTest"
  • :pulsar-broker:test --tests "org.apache.pulsar.broker.service.persistent.*"
  • :pulsar-broker:test --tests "org.apache.pulsar.broker.stats.ConsumerStatsTest" --tests "org.apache.pulsar.broker.stats.AuthenticatedConsumerStatsTest" --tests "org.apache.pulsar.broker.admin.*TopicStats*"
  • :pulsar-broker:test --tests MessageRedeliveryControllerTest --tests TransactionConsumeTest
  • :pulsar-common:test

Two places where coverage changes meaning rather than disappearing:

  • KeySharedSubscriptionTest.testContinueDispatchMessagesWhenMessageTTL had a && !impl.classic
    guard that suppressed the consumer3 assertion under the classic path. It is now unconditional,
    matching the already-unconditional consumer2 check above it — the assertion is strengthened, not
    dropped.
  • MessageRedeliveryControllerTest.testContainsStickyKeyHashes is rewritten onto the singular
    containsStickyKeyHash(int). The plural overload was an any-match loop, so
    containsStickyKeyHashes(Set.of(104, 105)) passed on 104 alone; the singular form now asserts 104
    present and 105 absent.

testRecentJoinedPosWillNotStuckOtherConsumer is kept — it guards
#7105 and is implementation-agnostic — with its
comments reworded into draining-hashes terms.

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

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

Public API: ConsumerStats.getReadPositionWhenJoining() and getKeyHashRanges() (both
@Deprecated since 4.0.0) and SubscriptionStats.getConsumersAfterMarkDeletePosition() are
removed, along with the corresponding ConsumerStatsImpl / SubscriptionStatsImpl fields.

Configuration: subscriptionKeySharedUseClassicPersistentImplementation and
subscriptionSharedUseClassicPersistentImplementation are removed. Both defaulted to false and
neither appeared in conf/broker.conf or conf/standalone.conf; a stale value is ignored rather
than rejected (see Upgrade impact).

REST endpoints: no endpoint changes; the topic-stats response shape loses
consumersAfterMarkDeletePosition (see Upgrade impact).

Documentation

  • doc-not-needed — the removed configuration keys were never documented in
    conf/broker.conf / conf/standalone.conf.

… in tests

### Motivation

The classic Shared/Key_Shared dispatchers are being removed in 5.0. Before the
implementation can go, the tests that parameterize on it have to stop doing so.
`KeySharedImplementationType` ran the four Key_Shared test classes twice via a
`@Factory`, starting two brokers for the largest of them.

### Modifications

- Delete `org.apache.pulsar.tests.KeySharedImplementationType` and unwind its
  four users: the `@Factory`, the extra constructors, the inject-only
  `currentImplementationType` data provider, the `prependImplementationTypeToData`
  wrapping of the real providers, and the classic config setters.
- `KeySharedSubscriptionTest.testContinueDispatchMessagesWhenMessageTTL`: the
  `!impl.classic` guards collapse to the modern form. Note the consumer3 check at
  the second site becomes unconditional, matching the already-unconditional
  consumer2 check above it.
- `testReadAheadLimit` and `testOrderingAfterReconnects` lose their
  `skipIfClassic()` and now always run.
- `KeySharedSubscriptionMaxUnackedMessagesTest` keeps its `KeySharedSelectorType`
  enum including `AutoSplit_Classic` — that constant selects the hash-range
  *selector* (`subscriptionKeySharedUseConsistentHashing`), an orthogonal axis.

No production code changes.
…assic implementation

### Motivation

Follow-up to the Key_Shared test cleanup, and the last step before the classic
Shared/Key_Shared dispatchers can be deleted. Four broker test classes still ran
every case twice, once against each dispatcher implementation.

### Modifications

- `ConsumerStatsTest`: `classicAndSubscriptionType` becomes
  `sharedSubscriptionTypes` (`Shared`, `Key_Shared`); `testConsumerStatsOutput`
  drops the `classicDispatchers` parameter and the config setters, and its
  expected-field set keeps only the modern `drainingHashes` /
  `keyHashRangeArrays` arm.
- `SharedDispatcherPermitAccountingTest`: `dispatcherImplementations` is gone and
  `flowRaceDispatcherVariants` collapses to the two subscription types;
  `createTestContext` loses its `classic` argument (the unused 1-arg overload is
  dropped) and the dispatcher branches, instance-of assertions and the two
  `totalAvailablePermits` helpers unwrap to the modern dispatcher.
  `untrackedSubscriptionVariants` is untouched — its first column is
  `persistent`, not `classic`.
- `SharedSubscriptionUnackedMessagesAccountingTest`: `@Factory`, the `classic`
  field and constructor go; one broker instead of two.
- `PersistentDispatcherMultipleConsumersTest`: `readConflationDispatcherTypes` is
  removed and `createReadConflationDispatcher` always builds the modern
  dispatcher.

No production code changes.
… implementations

### Motivation

PIP-379 replaced the Key_Shared "recently joined consumers" machinery with the
draining-hashes tracker. The implementation shipped in 4.0 LTS but kept the
pre-4.0 dispatchers alive behind two feature flags as a rollback escape hatch:

- `subscriptionKeySharedUseClassicPersistentImplementation`
- `subscriptionSharedUseClassicPersistentImplementation`

Both default to `false` and neither appears in `conf/broker.conf` or
`conf/standalone.conf`. Two LTS lines later the escape hatch has outlived its
purpose: ~2000 lines of duplicated dispatcher logic, a classic/modern branch
threaded through `PersistentSubscription`, and a path that lacks read-ahead
limiting while carrying the ordering defects (apache#23307, apache#21199) PIP-379 exists to
fix. 5.0 is the right major release to delete it.

### Modifications

- Delete `PersistentDispatcherMultipleConsumersClassic` and
  `PersistentStickyKeyDispatcherMultipleConsumersClassic`, their two dedicated
  test classes, and `NonEntryCacheKeySharedSubscriptionV30Test` (it forces the
  classic implementation and drives the classic-only
  `setSortRecentlyJoinedConsumersIfNeeded` / `getRecentlyJoinedConsumers` hooks).
- Delete both `ServiceConfiguration` flags.
- `PersistentSubscription`: the Shared branch collapses to
  `PersistentDispatcherMultipleConsumers`, and the Key_Shared branch loses only
  the classic arm — the PIP-486 `ksm.isEntryBucketDispatch()` arm stays.
- Drop the two commented-out `PULSAR_PREFIX_...` lines in
  `AbstractPulsarProfilingTest`.

A stale `subscriptionKeySharedUseClassicPersistentImplementation=true` left in
`broker.conf` is silently ignored — `FieldParser.update` iterates declared
fields, never the properties — and surfaces only in the `configOverrides`
attribute of the startup log. Both flags were `dynamic = true`, so a stale value
may sit in `/admin/configuration`; `BrokerService` already logs and skips unknown
dynamic keys.

The PIP-322 classic *dispatch rate limiter* is a different feature that first
shipped in 4.1.0 and is deliberately out of scope here.
…lementations used

### Motivation

With the classic dispatchers gone, the seams that existed solely to distinguish
them from the modern ones are dead weight.

### Modifications

- `isClassic()` is removed at all four sites: the abstract declaration on
  `AbstractPersistentDispatcherMultipleConsumers`, the `StickyKeyDispatcher`
  interface method, the `return false` on
  `PersistentDispatcherMultipleConsumers`, and the stats branch in
  `PersistentSubscription`, which now always emits `keyHashRangeArrays`. It was
  never called through the base type anywhere, including tests.
- `StickyKeyDispatcher.getRecentlyJoinedConsumers()` — a `default` returning
  `null` whose only overrider was the classic dispatcher — is removed together
  with the block in `PersistentSubscription` that fed
  `subStats.consumersAfterMarkDeletePosition` from it.
- `MessageRedeliveryController`: the `isClassicDispatcher` field and the 2-arg
  constructor are gone, `add()` loses its classic branches, and the write-only
  `containsStickyKeyHashes(Set<Integer>)` overload is removed in favour of the
  singular `containsStickyKeyHash(int)`.
- `MessageRedeliveryControllerTest`: drop
  `testClassicOutOfOrderDoesNotInitializePositionHashMap` and rewrite
  `testContainsStickyKeyHashes` onto the singular form. The plural assertions
  were an any-match loop, so `containsStickyKeyHashes(Set.of(104, 105))` passed
  on 104 alone; the singular version now asserts 104 present and 105 absent.
- Reword the stale "recent joined consumer tracking" comment in
  `PersistentStickyKeyDispatcherMultipleConsumers` — it describes what is
  actually the `drainingHashesRequired` flag.

`AbstractPersistentDispatcherMultipleConsumers` is deliberately kept: it is the
parameter type of the user-pluggable
`DelayedDeliveryTrackerFactory.newTracker(...)` SPI, the narrowing type
`BrokerService` uses for `blockedDispatchers`, and the type seven test files
mock.
…ats fields

### Motivation

With the classic dispatchers gone, three stats fields have no producer left.
This is the only commit in the series with a public admin-API break, so it is
kept as one reviewable diff.

- `ConsumerStats.getReadPositionWhenJoining()` and `getKeyHashRanges()` have been
  `@Deprecated` since apache#23429 (Oct 2024). They shipped deprecated in 4.0.0 and had
  4.0, 4.1 and 4.2 to warn.
- `SubscriptionStats.getConsumersAfterMarkDeletePosition()` was never marked
  `@Deprecated`, although `pip/pip-379.md` already declared it removed.

### Modifications

Remove all three from the admin API interfaces, the `*Impl` classes (field,
`add()`, `clear()`, getter) and `Consumer` (the `readPositionWhenJoining` field,
its `updateStats` branch, and `setReadPositionWhenJoining`, whose only caller was
the classic dispatcher).

### JSON-visible change

`consumersAfterMarkDeletePosition` was initialized to an empty `LinkedHashMap`,
so persistent subscriptions emit `"consumersAfterMarkDeletePosition": {}` today
and will stop. The other two are `null` on the modern path and the admin mapper
uses `Include.NON_NULL`, so their removal is JSON-invisible. The disappearing
`{}` is a deliberate 5.0 stats-shape change.

### Why `consumersAfterMarkDeletePosition` is removed, not renamed

It was never a general "consumers past the mark-delete position" statistic that
some other mechanism could repopulate — it was the classic dispatcher's
recently-joined-consumers map under a name that described its eviction rule:

- **Its own javadoc contradicts its name.** Both `SubscriptionStats` and
  `SubscriptionStatsImpl` carried the line *"This is for Key_Shared subscription
  to get the recentJoinedConsumers in the Key_Shared subscription."* The
  documented subject is `recentJoinedConsumers`.
- **The name describes a retention rule, not the data.** The key is a consumer
  name; the value is that consumer's cursor read position at the moment it
  joined — the same value surfaced as the consumer-level
  `readPositionWhenJoining`. "After mark delete position" referred only to the
  pruning invariant in `removeConsumersFromRecentJoinedConsumers()`, which
  dropped an entry once the mark-delete position passed the join position. The
  name conveys neither the key nor the value, only an incidental eviction
  condition.
- **It had no meaning outside the feature.** Its sole writer was
  `PersistentSubscription`, fed exclusively by
  `StickyKeyDispatcher.getRecentlyJoinedConsumers()`, whose only overrider was
  the classic dispatcher. PIP-282 renamed the consumer-level sibling
  (`readPositionWhenJoining` → `lastSentPositionWhenJoining`) but left this one
  alone, and PIP-379 then declared it removed outright.

`ConsumerStatsImpl` keeps its class-level `@SuppressWarnings("deprecation")` —
`lastAckedTimestamp` and `lastConsumedTimestamp` remain deprecated.
@lhotari lhotari added this to the 5.0.0 milestone Sep 22, 2026
@lhotari lhotari added release/important-notice The changes which are important should be mentioned in the release note and removed release/important-notice The changes which are important should be mentioned in the release note labels Sep 22, 2026

@void-ptr974 void-ptr974 left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. The cleanup is well scoped, the modern dispatcher path remains correct, and the test updates look solid. Nice work!

@lhotari
lhotari merged commit 9415e23 into apache:master Sep 22, 2026
44 checks passed
lhotari pushed a commit that referenced this pull request Oct 2, 2026
…ch index acks (#26809)

(cherry picked from commit 1046481)

[branch-4.x] Also applied to the classic Shared and Key_Shared dispatchers
(PersistentDispatcherMultipleConsumersClassic, PersistentStickyKeyDispatcherMultipleConsumersClassic),
which were removed on master by #26687. The test covers both implementations.
lhotari pushed a commit that referenced this pull request Oct 2, 2026
…ch index acks (#26809)

(cherry picked from commit 1046481)

[branch-4.x] Also applied to the classic Shared and Key_Shared dispatchers
(PersistentDispatcherMultipleConsumersClassic, PersistentStickyKeyDispatcherMultipleConsumersClassic),
which were removed on master by #26687. The test covers both implementations.
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.

4 participants