Skip to content

[improve][broker] Avoid the topic-wide deduplication lock - #26763

Merged
merlimat merged 1 commit into
apache:masterfrom
lhotari:lh-perfopt-dedup-producer-state
Sep 29, 2026
Merged

merlimat merged 1 commit into
apache:masterfrom
lhotari:lh-perfopt-dedup-producer-state

Conversation

@lhotari

@lhotari lhotari commented Sep 29, 2026 •

Copy link
Copy Markdown
Member

Motivation

Broker-side duplicate admission synchronizes on one map for the entire topic (synchronized (highestSequencedPushed) in MessageDeduplication.isDuplicateNormal). Messages from different producers therefore contend on the same monitor before their managed-ledger adds are submitted. This becomes expensive when many producer connections publish to one topic, even though sequence ordering only needs coordination between calls that use the same producer name.

In the iot-telemetry-max-rate performance scenario (500 producers publishing unbatched messages to one topic with deduplication enabled), the pulsar-io threads blocked 3.0 s on this monitor during a profiled 8 KB run on master; with this change that wait is gone.

Modifications

Use the existing ConcurrentMap operations to update the highest pushed sequence ID with a putIfAbsent/conditional replace loop. Updates for one producer remain linearizable, including the brief overlap possible during a reconnect, while unrelated producers no longer serialize on a topic-wide monitor.

The replication-v2 path retains its existing synchronization because its ledger-ID and entry-ID values must be checked and updated as one pair. Snapshot, recovery, persistence, and inactive-producer cleanup retain the existing map representation and behavior.

Measurements

Scenario iot-telemetry-max-rate of tests/performance (500 producers, one topic, one application with a 20-member Key_Shared subscription, unbatched messages, deduplication enabled), run with the performance launcher on one host (Intel i9-9980HK, 8 cores, 16 hardware threads, turbo off at a fixed 2.4 GHz, thermald stopped, cool-down to 55 °C before each run). Baseline is master f0c9094aa3fa, candidate this change 9e11f9d0a3e0 on top of it, each in its own image. Both sides ran with the bookies' DbLedgerStorage caches at 256 MB (--set cluster.bookies.env.dbStorage_writeCacheMaxSizeMb=256 and dbStorage_readAheadCacheMaxSizeMb=256), because the test image's 16 MB default throttles bookie writes of large entries; #26749 makes that the scenarios' default. The unprofiled runs were interleaved (baseline, change, change, baseline, baseline, change) for each payload size. Every run delivered every message with no duplicates, ordering violations or invalid messages, and none throttled.

Unprofiled, means of 3 Master This change Change
128 B throughput 109,043 msg/s (107.2k–110.7k) 108,451 msg/s (107.0k–110.7k) −0.5 %
128 B publish / end-to-end p99 1,050 / 1,060 ms 1,080 / 1,086 ms +3 %, within spread
8 KB throughput 43,564 msg/s (37.5k–47.6k) 46,729 msg/s (44.2k–48.7k) +7.3 %
8 KB publish / end-to-end p99 1,048 / 1,062 ms 374 / 385 ms −64 %
  • 128 B: no difference. On master, this size is limited by the topic's managed-ledger thread, which is saturated, not by the deduplication monitor.
  • 8 KB: the +7.3 % and the −64 % p99 both come from one slow master run (37.5k msg/s with a 2.3 s p99). The other two master runs (45.6k and 47.6k msg/s) are within this change's range, so the throughput difference is inconclusive with three runs each.
  • Profiled 8 KB pair (--extends configs/profile-broker): throughput is the same (49.7k vs 49.1k msg/s). The blocked time on the deduplication monitor drops from 3.0 s to none. The largest blocked time on both sides is then the dispatcher's Flow monitor (12–13 s), which [improve][broker] Apply Shared and Key_Shared Flow permits without the dispatcher monitor #26716 removes.

This change removes the monitor waits, but on master that monitor does not limit this scenario's throughput. In an experiment branch with #26716 and add submission moved off the managed-ledger thread, the deduplication monitor had become the largest blocked wait at 8 KB (5.1 s of 9.3 s blocked with an application frame), and it disappeared with this change.

Charts of the median unprofiled 8 KB run of each side (master 45,613 msg/s, this change 47,302 msg/s), as the run report renders them:

Note

Each chart has its own scale, set by that run's values, so the baseline's and this change's charts of the same measure have different axes. Compare their values, not the heights of the lines.

Latency by percentile, baseline (master)

Latency by percentile of the baseline

Latency by percentile, this change

Latency by percentile with this change

Throughput over time, baseline (master)

Throughput over time of the baseline

Throughput over time, this change

Throughput over time with this change

Backlog over time, baseline (master)

Backlog over time of the baseline

Backlog over time, this change

Backlog over time with this change

Earlier component measurements with a focused JMH benchmark (MessageDeduplicationSequenceCheckBenchmark):

Case Topic-wide monitor Per-producer update Change
Single producer thread 41.1 ns/op 31.6 ns/op 23 % faster
16 distinct producers 1,435.8 ns/op 114.3 ns/op 92 % faster
16 callers sharing one producer 1,077.1 ns/op 1,743.7 ns/op 62 % slower

The shared-producer case models an exceptional reconnect overlap. Normal operation has one active connection for a producer name, while distinct producer names are the case this change improves. Accepted updates allocate 24 bytes in both implementations; retries in the deliberately contended same-producer case raised the candidate to about 31 bytes/op.

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • MessageDuplicationTest covers a 32-thread same-producer admission race; it and MessageDeduplicationTest pass on this branch (11 tests).
  • MessageDeduplicationSequenceCheckBenchmark in microbench measures the admission path with distinct and shared producers.
  • ./gradlew spotlessCheck checkstyleMain checkstyleTest for the changed modules passes.

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

If the box was checked, please highlight the changes

  • 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

This PR was prepared with AI assistance (Claude Code) and reviewed by a human contributor.

@merlimat
merlimat merged commit 16e7ee0 into apache:master Sep 29, 2026
45 checks passed
@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