Repository navigation
[improve][broker] Avoid the topic-wide deduplication lock - #26763
Merged
Merged
Conversation
Assisted-by: Codex
merlimat
approved these changes
Sep 29, 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.
Motivation
Broker-side duplicate admission synchronizes on one map for the entire topic (
synchronized (highestSequencedPushed)inMessageDeduplication.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-rateperformance 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
ConcurrentMapoperations to update the highest pushed sequence ID with aputIfAbsent/conditionalreplaceloop. 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-rateoftests/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,thermaldstopped, cool-down to 55 °C before each run). Baseline is masterf0c9094aa3fa, candidate this change9e11f9d0a3e0on 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=256anddbStorage_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.--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, this change
Throughput over time, baseline (master)
Throughput over time, this change
Backlog over time, baseline (master)
Backlog over time, this change
Earlier component measurements with a focused JMH benchmark (
MessageDeduplicationSequenceCheckBenchmark):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
This change added tests and can be verified as follows:
MessageDuplicationTestcovers a 32-thread same-producer admission race; it andMessageDeduplicationTestpass on this branch (11 tests).MessageDeduplicationSequenceCheckBenchmarkinmicrobenchmeasures the admission path with distinct and shared producers../gradlew spotlessCheck checkstyleMain checkstyleTestfor the changed modules passes.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes
This PR was prepared with AI assistance (Claude Code) and reviewed by a human contributor.