Skip to content

[improve][cli] pulsar-perf: cut latency histogram memory ~100x, fix range units and static stats state - #26466

Merged
merlimat merged 2 commits into
apache:masterfrom
lhotari:lh-improve-perf-histogram-digits
Sep 5, 2026
Merged

merlimat merged 2 commits into
apache:masterfrom
lhotari:lh-improve-perf-histogram-digits

Conversation

@lhotari

@lhotari lhotari commented Sep 4, 2026 •

Copy link
Copy Markdown
Member

Motivation

A heap dump of a pulsar-perf produce run (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:

class PerformanceProducer     88,090,448   8.52%
class PerformanceTransaction  88,084,968   8.52%
class PerformanceConsumer     29,363,080   2.84%
class PerformanceReader       29,362,392   2.84%
class ManagedLedgerWriter     23,073,760   2.23%

Only produce was running. Investigating this turned up two independent problems.

1. The histograms are built at HdrHistogram's maximum precision

All 13 Recorder construction sites in pulsar-testclient used numberOfSignificantValueDigits = 5, which is the maximum the API allows. HdrHistogram sizes a fixed-range histogram's counts array as (bucketsNeeded + 1) * subBucketHalfCount, where subBucketHalfCount = 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 range PerformanceProducer used:

significant digits counts array, per Recorder
3 229,888 B
4 3,146,240 B
5 22,020,608 B

Because these recorders are static fields and PulsarPerfTestTool.initCommander reflectively instantiates every registered subcommand before parsing arguments (PulsarPerfTestTool.java:72, constructor.newInstance()), each class's <clinit> runs on every invocation. So pulsar-perf produce was also paying for transaction, consume, read and managed-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%:

  p50     5-digit=    2.001 ms   3-digit=    2.001 ms   rel.err=0.0000%
  p99     5-digit=   12.840 ms   3-digit=   12.847 ms   rel.err=0.0545%
  p99.99  5-digit=   38.243 ms   3-digit=   38.271 ms   rel.err=0.0732%
  p100    5-digit=  143.872 ms   3-digit=  143.999 ms   rel.err=0.0883%

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 Recorder is not the fix: Recorder(int) uses ConcurrentHistogram, 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 throws ArrayIndexOutOfBoundsException on an out-of-range value. #5786 fixed such a mismatch in PerformanceProducer by switching its range from SECONDS.toMillis(120000) to SECONDS.toMicros(120000), but the same bug was left in place elsewhere:

  • ManagedLedgerWriter records latencyMicros (NANOSECONDS.toMicros(...)) into a range of SECONDS.toMillis(120000) = 1.2e8. Its real ceiling is therefore about two minutes, not 120000 seconds. There is no clamp and no handler — the recordValue sits directly in the AddEntryCallback.addComplete callback. Present since the file was added in Added ManagedLedger perf tool #1270.
  • SimpleTestProducerSocket (pulsar-perf websocket-producer) has exactly the same mismatch, inside an @OnWebSocketMessage handler. Dates to the initial import.
  • PerformanceReader records a wall-clock latencyMillis with only a >= 0 guard. PerformanceConsumer was 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"); PerformanceReader never was, so reading a sufficiently old backlog throws out of the read loop.

Note that the declared highestTrackableValue is 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 static at all, and neither do the 51 LongAdder counters sitting next to them in the same declaration blocks. The only static accessors are printAggregatedStats and printAggregatedThroughput, and those are called from a shutdown-hook lambda that is already constructed inside run():

// PerformanceReader.java:156 — inside run(), the lambda already captures `this`
Thread shutdownHookThread = addShutdownHook(() -> {
    printAggregatedThroughput(start);
    printAggregatedStats();
});

PerformanceProducer and PerformanceConsumer had in fact already made printAggregatedThroughput an instance method, so their counters had no static accessor left.

What the static did cost is 26 defensive reset() calls, under a comment that names the problem:

public void run() throws Exception {
    // Reset static counters to avoid stale state from previous runs in the same JVM
    messagesSent.reset();
    ...

and the compensation is incomplete: PerformanceReader, ManagedLedgerWriter and PerformanceClient have 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. PerformanceProducerTest already constructs four PerformanceProducer instances 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

  • Add PerfClientUtils.LATENCY_HISTOGRAM_SIGNIFICANT_DIGITS = 3 and use it at all 13 Recorder construction sites, with the sizing rationale documented on the constant.
  • Give each publish-side recorder a MAX_LATENCY_MICROS = TimeUnit.HOURS.toMicros(1) ceiling (PerformanceProducer, PerformanceTransaction, ManagedLedgerWriter, SimpleTestProducerSocket) — correcting the unit for the latter two — and clamp with Math.min before recording. Consume/read keep TimeUnit.DAYS.toMillis(10).
  • Clamp in PerformanceReader, matching what PerformanceConsumer already does.
  • Rename PerformanceConsumer.MAX_LATENCY to MAX_LATENCY_MILLIS; unit-ambiguous names are what allowed this class of bug in the first place.
  • Drop the now-unnecessary ArrayIndexOutOfBoundsException branch in PerformanceProducer's exceptionally handler. It was a workaround for the recorder throwing; with clamping it can only mask a genuine error.

And, for the JVM-global state:

  • Make all 12 Recorder and 51 LongAdder fields instance fields, and printAggregatedStats / printAggregatedThroughput / printTxnAggregatedThroughput instance methods, across PerformanceProducer, PerformanceConsumer, PerformanceReader, PerformanceTransaction, ManagedLedgerWriter and PerformanceClient.
  • Delete the 26 compensating reset() calls; the state now cannot outlive the command object.
  • SimpleTestProducerSocket.recorder was the one field with a genuine reason to be shared, since PerformanceClient creates one socket per connection. PerformanceClient now owns the Recorder and 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 produce run (the 14 arrays MAT observed):

  before = 257,956,864 bytes (246.0 MB)
  after  =   2,579,456 bytes (2.46 MB)   -> 100x smaller

Reported percentiles shift by at most ~0.09% (visible only in the third decimal of the millisecond figures); the %d millisecond 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-perf exec-ed one main class per subcommand (exec $JAVA $OPTS org.apache.pulsar.testclient.PerformanceProducer ...), so only the selected command's <clinit> ran. Registering the subcommand Class objects with picocli instead of pre-built instances would restore that property — picocli defers CommandUserObject.getInstance() until a subcommand is actually selected.

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • The de-staticizing is verified by the compiler: had any of the 63 fields genuinely needed to be static, the build would fail. PerformanceProducerTest, which constructs four PerformanceProducer instances in one JVM, passes without the reset calls that previously made that safe.
  • Added PerfClientUtilsTest.latencyHistogramsStaySmallAtTheConfiguredPrecision, which asserts that a Recorder built at LATENCY_HISTOGRAM_SIGNIFICANT_DIGITS over 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.
  • The array-sizing figures above were measured by running against HdrHistogram-2.2.2.jar directly, 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.
  • Also confirmed that this is not a regression from the recent HdrHistogram 2.1.9 -> 2.2.2 bump ([improve][monitor] Upgrade Dropwizard Metrics to 4.2.39 and HdrHistogram to 2.2.2 #26353). Both versions compute byte-identical counts-array sizes for these declarations. If anything the bump helped: in 2.1.9 Recorder allocated both the active and the inactive histogram up front (measured 44,040,224 B per Recorder at construction), whereas 2.2.2 leaves inactiveHistogram null until the first getIntervalHistogram() (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

  • 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

…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)
…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)
@lhotari lhotari changed the title [improve][cli] pulsar-perf: cut latency histogram memory ~100x and fix range unit mismatches [improve][cli] pulsar-perf: cut latency histogram memory ~100x, fix range units and static stats state Sep 4, 2026
@lhotari

lhotari commented Sep 4, 2026

Copy link
Copy Markdown
Member Author

Follow-up: #26467 fixes the structural half of this — PulsarPerfTestTool eagerly instantiating every subcommand, which is what makes produce pay for consume/read/transaction/managed-ledger in the first place. The two are independent and can land in either order.

@merlimat
merlimat merged commit bc1785d into apache:master Sep 5, 2026
44 checks passed
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.
@david-streamlio
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.
@lhotari lhotari added this to the 5.0.0-M2 milestone Sep 9, 2026
lhotari added a commit that referenced this pull request Sep 9, 2026
…ange units and static stats state (#26466)

(cherry picked from commit bc1785d)
lhotari added a commit that referenced this pull request Sep 9, 2026
…ange units and static stats state (#26466)

(cherry picked from commit bc1785d)
Radiancebobo pushed a commit to Radiancebobo/pulsar that referenced this pull request Oct 8, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants