Repository navigation
[improve][cli] pulsar-perf: cut latency histogram memory ~100x, fix range units and static stats state - #26466
Merged
Conversation
…x range unit mismatches A heap dump of a `pulsar-perf produce` run showed 246 MB — 25% of the client heap — held by HdrHistogram counts arrays before a single message was sent. All 13 Recorder construction sites used numberOfSignificantValueDigits = 5, HdrHistogram's maximum. The counts array grows by roughly 10x per digit, so a single Recorder over the ranges used here costs 11-22 MB. Because the recorders are static fields and PulsarPerfTestTool reflectively instantiates every subcommand at startup, `produce` also paid for `consume`, `read`, `transaction` and `managed-ledger`. Nothing reported resolves that precision: the microsecond-based tools print milliseconds with %.3f and the millisecond-based tools print whole milliseconds with %d. ManagedLedgerWriter and SimpleTestProducerSocket also had the unit mismatch that apache#5786 fixed in PerformanceProducer but never propagated: both record microseconds into a range sized with SECONDS.toMillis(120000), giving a real ceiling of 120 s rather than 120000 s, with no clamp and no handler. PerformanceReader never got the clamp PerformanceConsumer received in apache#17160, so reading a backlog older than the histogram range throws out of the read loop. Assisted-by: Claude Opus 5 (Claude Code)
lhotari
requested review from
Technoboy-,
dao-jun,
david-streamlio,
merlimat and
nodece
September 4, 2026 20:19
…tance fields The latency Recorders and the 51 LongAdder counters in the perf commands were static for no reason: the only static accessors were printAggregatedStats and printAggregatedThroughput, which are called from a shutdown-hook lambda already constructed inside run(). PerformanceProducer and PerformanceConsumer had already made printAggregatedThroughput an instance method, leaving their counters with no static accessor at all. The cost of that was 26 defensive reset() calls under the comment "Reset static counters to avoid stale state from previous runs in the same JVM" — and the compensation was incomplete: PerformanceReader, ManagedLedgerWriter and PerformanceClient have no resets, so a second in-process run folded the first run's samples into its aggregate stats. PerformanceProducerTest already constructs four PerformanceProducer instances in one JVM. SimpleTestProducerSocket.recorder was the one field with a real reason to be shared, since PerformanceClient creates a socket per connection. It is now owned by PerformanceClient and injected, which keeps the aggregation while dropping the JVM-global. Assisted-by: Claude Opus 5 (Claude Code)
1 of 11 tasks
Member
Author
|
Follow-up: #26467 fixes the structural half of this — |
Merged
1 of 11 tasks
merlimat
approved these changes
Sep 5, 2026
lhotari
added a commit
to lhotari/pulsar
that referenced
this pull request
Sep 7, 2026
…5 ones ### Motivation `pulsar-perf` was migrated to the V5 client API in apache#25887 (and apache#25917 did the same for `pulsar-client`), so every perf subcommand now drives the V5 SDK. That leaves no way to benchmark the v4 client, or v4 (non-scalable) topics, without the V5 SDK in the path — and several v4 capabilities became unreachable, because they have no V5 equivalent: - `--max-outstanding` / `--max-outstanding-across-partitions` and round-robin partition routing on the producer; - real `Exclusive` / `Failover` / `Key_Shared` subscription types, `MessageListener` dispatch on the client's listener threads, pooled messages, batch-index acknowledgment, the chunked-message knobs, the receiver-queue limits and the auto-scaled receiver-queue reporting on the consumer; - the v4 `Reader` itself — `read` measures a V5 `CheckpointConsumer`, a different broker-side entity — along with a `lid:eid` start message id, `--receiver-queue-size` and `--use-tls`; - the v4 transaction coordinator, which stays a live broker code path next to the v5 one (PIP-473 P5.4) with nothing driving it, and its acknowledgement round-trip latency (V5's `acknowledge` is a synchronous void, so its reported ack latency is a local measurement). ### Modifications Each benchmark is split into an abstract base holding everything that is not client-specific — the CLI options, the throughput/latency accounting, the run() skeleton and message loop, and the reports — plus two thin subclasses that bind the client types: the existing V5 command, and a new v4 one. PerformanceProducerBase -> PerformanceProducer / PerformanceProducerV4 (produce-v4) PerformanceConsumerBase -> PerformanceConsumer / PerformanceConsumerV4 (consume-v4) PerformanceReaderBase -> PerformanceReader / PerformanceReaderV4 (read-v4) PerformanceTransactionBase -> PerformanceTransaction / PerformanceTransactionV4 (transaction-v4) This is a split rather than a revert, so the bases keep the improvements made since the migration — instance recorders, the 3-significant-digit histograms and the latency clamps from apache#26466, class-based subcommand registration from apache#26467 — and both commands share one benchmark and one set of measurements. The v4 commands offer the same flags as their V5 counterparts (only the V5-only `--scalable-consumer-type` and `--scalable` are not mirrored). Each base logs through a logger named after the concrete subclass, so the report lines keep naming the subcommand that produced them. Where the two clients genuinely differ, the seam is explicit: the v4 consumer and reader keep `MessageListener`/`ReaderListener` dispatch while V5 keeps its poll threads; the v4 producer does not await a transaction's sends before committing (V5 must, because its transactional sends are queued onto an internal dispatch chain); only the first transaction of a V5 test thread waits out the coordinator's asynchronous connect, so the rollover loop still counts every failed open. `--jsse-provider` / `--jca-provider` are also wired into the v4 client builder, which previously read them only for the admin and V5 legs. ### Verifying this change New unit and end-to-end tests in `pulsar-testclient` cover subcommand registration and naming, conf-file defaults and case-insensitive enums on the v4 commands, the shared `-st` enum mapping onto the v4 client enum, produce / consume / read round trips through the v4 commands, the `lid:eid` start position, and the v4-only producer knobs. `PerfToolTest` gains v4 cases that run the commands from `bin/pulsar-perf` in a container.
lhotari
added a commit
to lhotari/pulsar
that referenced
this pull request
Sep 7, 2026
…5 ones ### Motivation `pulsar-perf` was migrated to the V5 client API in apache#25887 (and apache#25917 did the same for `pulsar-client`), so every perf subcommand now drives the V5 SDK. That leaves no way to benchmark the v4 client, or v4 (non-scalable) topics, without the V5 SDK in the path — and several v4 capabilities became unreachable, because they have no V5 equivalent: - `--max-outstanding` / `--max-outstanding-across-partitions` and round-robin partition routing on the producer; - real `Exclusive` / `Failover` / `Key_Shared` subscription types, `MessageListener` dispatch on the client's listener threads, pooled messages, batch-index acknowledgment, the chunked-message knobs, the receiver-queue limits and the auto-scaled receiver-queue reporting on the consumer; - the v4 `Reader` itself — `read` measures a V5 `CheckpointConsumer`, a different broker-side entity — along with a `lid:eid` start message id, `--receiver-queue-size` and `--use-tls`; - the v4 transaction coordinator, which stays a live broker code path next to the v5 one (PIP-473 P5.4) with nothing driving it, and its acknowledgement round-trip latency (V5's `acknowledge` is a synchronous void, so its reported ack latency is a local measurement). ### Modifications Each benchmark is split into an abstract base holding everything that is not client-specific — the CLI options, the throughput/latency accounting, the run() skeleton and message loop, and the reports — plus two thin subclasses that bind the client types: the existing V5 command, and a new v4 one. PerformanceProducerBase -> PerformanceProducer / PerformanceProducerV4 (produce-v4) PerformanceConsumerBase -> PerformanceConsumer / PerformanceConsumerV4 (consume-v4) PerformanceReaderBase -> PerformanceReader / PerformanceReaderV4 (read-v4) PerformanceTransactionBase -> PerformanceTransaction / PerformanceTransactionV4 (transaction-v4) This is a split rather than a revert, so the bases keep the improvements made since the migration — instance recorders, the 3-significant-digit histograms and the latency clamps from apache#26466, class-based subcommand registration from apache#26467 — and both commands share one benchmark and one set of measurements. The v4 commands offer the same flags as their V5 counterparts (only the V5-only `--scalable-consumer-type` and `--scalable` are not mirrored). Each base logs through a logger named after the concrete subclass, so the report lines keep naming the subcommand that produced them. Where the two clients genuinely differ, the seam is explicit: the v4 consumer and reader keep `MessageListener`/`ReaderListener` dispatch while V5 keeps its poll threads; the v4 producer does not await a transaction's sends before committing (V5 must, because its transactional sends are queued onto an internal dispatch chain); only the first transaction of a V5 test thread waits out the coordinator's asynchronous connect, so the rollover loop still counts every failed open. `--jsse-provider` / `--jca-provider` are also wired into the v4 client builder, which previously read them only for the admin and V5 legs. ### Verifying this change New unit and end-to-end tests in `pulsar-testclient` cover subcommand registration and naming, conf-file defaults and case-insensitive enums on the v4 commands, the shared `-st` enum mapping onto the v4 client enum, produce / consume / read round trips through the v4 commands, the `lid:eid` start position, and the v4-only producer knobs. `PerfToolTest` gains v4 cases that run the commands from `bin/pulsar-perf` in a container.
2 of 14 tasks
david-streamlio
removed their request for review
September 7, 2026 19:27
lhotari
added a commit
to lhotari/pulsar
that referenced
this pull request
Sep 7, 2026
…ands Follow-up to the review on apache#26480, pulsar-perf half. - `transaction-v4` had no functional coverage: every existing reference was metadata-only, so the transaction binding itself was unpinned. Adds `PerformanceTransactionV4Test` with a commit run and an `-abort` run against plain `persistent://` topics. The commit run asserts both halves are durable (produced messages visible, consume backlog down by one per transaction); the abort run asserts neither is (nothing visible, backlog unchanged, and all ten messages redeliverable on the tool's own subscription, which is what separates a real abort from a transaction that was merely never ended). Mutation checked: making `sendMessage` or `acknowledgeAsync` ignore the transaction, dropping the acknowledgement entirely, or making `abortTransaction` a no-op each fail it. - `PerformanceConsumerBase` had dropped the `consumerType` attribute from the per-topic "Adding consumers" line, which the pre-split `PerformanceConsumer` logged. Restores it through a `consumerTypeForLog()` hook: the base reports the shared `--subscription-type`, `PerformanceConsumer` overrides it with `scalableConsumerType`, so the V5 line matches the pre-split output again. - `PerformanceTransactionBase` interrupted `Thread.currentThread()` from its failure callbacks. That is the worker thread for a V5 ack, whose future is already complete when returned, but a client-internal thread for a v4 ack and for the send and end-transaction futures on both clients. All three callbacks now capture the worker thread, which is the thread `runWorker` invokes them on. - `transaction-v4`'s and `produce-v4`'s `sendTimeout(0)` was annotated "a send timeout and a transaction are mutually exclusive on the v4 client". That guard was removed in apache#16519; the setting is still right because a send timeout fails the send on its own schedule and takes the transaction with it, so the comments now say that instead. - `PerformanceV4CommandsTest`: the backlog assertion now waits with Awaitility (the exit latch is released before the client close that flushes the acks, and the 30s join is unchecked), and the delivery-time assertions also pin each message's payload so a future reordering fails legibly. - The `LATENCY_HISTOGRAM_SIGNIFICANT_DIGITS` rationale still described the static recorders that apache#26466 replaced with instance fields in the same commit that added the comment. Rewords both copies, with the figure re-measured on HdrHistogram 2.2.2: 16.00 MB (1h in micros) and 14.00 MB (10d in millis) at 5 digits, not the pre-apache#26466 ranges' 11-22 MB.
Radiancebobo
pushed a commit
to Radiancebobo/pulsar
that referenced
this pull request
Oct 8, 2026
…ange units and static stats state (apache#26466)
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
A heap dump of a
pulsar-perf producerun (1.03 GB live) showed 246 MB — 25% of the client heap — held by HdrHistogram counts arrays before a single message had been sent. Eclipse MAT attributed it to class objects, i.e. static fields:Only
producewas running. Investigating this turned up two independent problems.1. The histograms are built at HdrHistogram's maximum precision
All 13
Recorderconstruction sites inpulsar-testclientusednumberOfSignificantValueDigits = 5, which is the maximum the API allows. HdrHistogram sizes a fixed-range histogram's counts array as(bucketsNeeded + 1) * subBucketHalfCount, wheresubBucketHalfCount = 2^ceil(log2(2 * 10^digits)) / 2— so each additional digit multiplies the allocation by roughly 10. Measured against HdrHistogram 2.2.2 over the rangePerformanceProducerused:RecorderBecause these recorders are
staticfields andPulsarPerfTestTool.initCommanderreflectively instantiates every registered subcommand before parsing arguments (PulsarPerfTestTool.java:72,constructor.newInstance()), each class's<clinit>runs on every invocation. Sopulsar-perf producewas also paying fortransaction,consume,readandmanaged-ledger— about 170 MB of histograms for commands that never ran.The 5th digit is not visible anywhere in the output. The microsecond-based tools divide by 1000 and print milliseconds with
%7.3f; the millisecond-based tools print whole milliseconds with%d. Over 2M synthetic lognormal samples, 3 vs 5 digits differ by at most 0.088%:That is far below the run-to-run variance of a throughput benchmark, and 3 digits is HdrHistogram's own recommended default. The quantization is one-sided upward for percentiles and Max, so reduced precision can only inflate a reported tail, never understate it.
The precision change is also provably incapable of introducing a new
ArrayIndexOutOfBoundsException: the throw threshold is a function of the range alone. Measured first-throwing values are byte-identical at 5 and at 3 digits for every range used here.Note that switching to an auto-resizing
Recorderis not the fix:Recorder(int)usesConcurrentHistogram, which keeps both an active and an inactive counts array. Measured,new Recorder(5)starts at 4.2 MB and grows to 44 MB — twice the fixed-range cost.2. Two ranges are in the wrong unit, and two record sites can throw
new Recorder(highestTrackableValue, digits)builds a fixed-range histogram that throwsArrayIndexOutOfBoundsExceptionon an out-of-range value. #5786 fixed such a mismatch inPerformanceProducerby switching its range fromSECONDS.toMillis(120000)toSECONDS.toMicros(120000), but the same bug was left in place elsewhere:ManagedLedgerWriterrecordslatencyMicros(NANOSECONDS.toMicros(...)) into a range ofSECONDS.toMillis(120000)= 1.2e8. Its real ceiling is therefore about two minutes, not 120000 seconds. There is no clamp and no handler — therecordValuesits directly in theAddEntryCallback.addCompletecallback. Present since the file was added in Added ManagedLedger perf tool #1270.SimpleTestProducerSocket(pulsar-perf websocket-producer) has exactly the same mismatch, inside an@OnWebSocketMessagehandler. Dates to the initial import.PerformanceReaderrecords a wall-clocklatencyMilliswith only a>= 0guard.PerformanceConsumerwas given an upper clamp in [bug] pulsar-perf consume, do not fail in case of reading data older than 10 days #17160 ("do not fail in case of reading data older than 10 days");PerformanceReadernever was, so reading a sufficiently old backlog throws out of the read loop.Note that the declared
highestTrackableValueis not itself the throw threshold — HdrHistogram rounds the counts array up to the next power-of-two bucket boundary, so each histogram tolerates 12-24% more than its declared range. Measured first-throwing values: 134,217,728 us (134.2 s) for the two mismatched ranges, and 1,073,741,824 ms (12.43 days) for consume/read. The bugs stand; only the exact ceilings differ from the declared ones.Separately, the 33.3-hour (
SECONDS.toMicros(120000)) ceiling on the publish-side recorders is not a meaningful latency bound for a send or an ack. Consume and read keep their 10-day range on purpose, because a backlogged message legitimately can be that old.3. The same fields are JVM-global for no reason
While fixing the above it became clear the recorders did not need to be
staticat all, and neither do the 51LongAddercounters sitting next to them in the same declaration blocks. The only static accessors areprintAggregatedStatsandprintAggregatedThroughput, and those are called from a shutdown-hook lambda that is already constructed insiderun():PerformanceProducerandPerformanceConsumerhad in fact already madeprintAggregatedThroughputan instance method, so their counters had no static accessor left.What the
staticdid cost is 26 defensivereset()calls, under a comment that names the problem:and the compensation is incomplete:
PerformanceReader,ManagedLedgerWriterandPerformanceClienthave no resets at all, so a second run in the same JVM folds the first run's messages, bytes and latencies into its aggregate stats.PerformanceProducerTestalready constructs fourPerformanceProducerinstances in one JVM.This is handled here rather than in a follow-up because it is the same fields, in the same declaration blocks, in the same six classes this PR already rewrites — splitting it would mean touching every one of those lines twice for no reviewer benefit.
Modifications
PerfClientUtils.LATENCY_HISTOGRAM_SIGNIFICANT_DIGITS = 3and use it at all 13Recorderconstruction sites, with the sizing rationale documented on the constant.MAX_LATENCY_MICROS = TimeUnit.HOURS.toMicros(1)ceiling (PerformanceProducer,PerformanceTransaction,ManagedLedgerWriter,SimpleTestProducerSocket) — correcting the unit for the latter two — and clamp withMath.minbefore recording. Consume/read keepTimeUnit.DAYS.toMillis(10).PerformanceReader, matching whatPerformanceConsumeralready does.PerformanceConsumer.MAX_LATENCYtoMAX_LATENCY_MILLIS; unit-ambiguous names are what allowed this class of bug in the first place.ArrayIndexOutOfBoundsExceptionbranch inPerformanceProducer'sexceptionallyhandler. It was a workaround for the recorder throwing; with clamping it can only mask a genuine error.And, for the JVM-global state:
Recorderand 51LongAdderfields instance fields, andprintAggregatedStats/printAggregatedThroughput/printTxnAggregatedThroughputinstance methods, acrossPerformanceProducer,PerformanceConsumer,PerformanceReader,PerformanceTransaction,ManagedLedgerWriterandPerformanceClient.reset()calls; the state now cannot outlive the command object.SimpleTestProducerSocket.recorderwas the one field with a genuine reason to be shared, sincePerformanceClientcreates one socket per connection.PerformanceClientnow owns theRecorderand injects it through the constructor, which keeps the aggregation while dropping the JVM-global.That part is a pure refactor with no behaviour change for a single run, and it is net-negative in size (-116 / +81 lines). One piece of mutable static state remains in
PerformanceClient(messageFormatter); it is unrelated to the stats and left alone.Measured effect on the startup allocation for a
pulsar-perf producerun (the 14 arrays MAT observed):Reported percentiles shift by at most ~0.09% (visible only in the third decimal of the millisecond figures); the
%dmillisecond reports are unchanged for any value below ~2 seconds, where 3-digit buckets are still exact.This does not address the underlying issue that every subcommand is instantiated at startup. That is worth fixing separately and would recover the ~170 MB attributable to commands that never run. It is also a comparatively recent regression: before #22388 (PIP-343, May 2024)
bin/pulsar-perfexec-ed one main class per subcommand (exec $JAVA $OPTS org.apache.pulsar.testclient.PerformanceProducer ...), so only the selected command's<clinit>ran. Registering the subcommandClassobjects with picocli instead of pre-built instances would restore that property — picocli defersCommandUserObject.getInstance()until a subcommand is actually selected.Verifying this change
This change added tests and can be verified as follows:
static, the build would fail.PerformanceProducerTest, which constructs fourPerformanceProducerinstances in one JVM, passes without the reset calls that previously made that safe.PerfClientUtilsTest.latencyHistogramsStaySmallAtTheConfiguredPrecision, which asserts that aRecorderbuilt atLATENCY_HISTOGRAM_SIGNIFICANT_DIGITSover each range the perf clients use stays under 512 KB. Verified that it pins the change: it fails with the constant set back to 5 (Expecting actual ... to be less than 524288) and passes at 3.HdrHistogram-2.2.2.jardirectly, and the three array lengths seen in the heap dump (2,752,512 / 1,835,008 / 1,441,792) reproduce exactly from the declared ranges.Recorderallocated both the active and the inactive histogram up front (measured 44,040,224 B perRecorderat construction), whereas 2.2.2 leavesinactiveHistogramnull until the firstgetIntervalHistogram()(22,020,112 B). The same dump taken before [improve][monitor] Upgrade Dropwizard Metrics to 4.2.39 and HdrHistogram to 2.2.2 #26353 would have carried 24 arrays / 427,819,392 B rather than 14 arrays / 257,949,920 B. The ranges and the digit count are the Pulsar-side cause, and they predate both versions.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes