Repository navigation
[improve][test] Migrate pulsar-perf to the V5 client API - #25887
Merged
Merged
Conversation
PulsarClientBuilderV5#tlsPolicy was a stub that only set useTls=true and dropped every field on the TlsPolicy on the floor. Wire the policy through: trustCertsFilePath, keyFilePath, certificateFilePath, allowInsecureConnection, enableHostnameVerification all land on the underlying ClientConfigurationData. This unblocks any V5 caller that needs custom TLS — most immediately the pulsar-perf migration where users pass --trust-cert-file and expect the broker connection to honor it. Added two unit tests covering field propagation and the TlsPolicy.ofInsecure() shortcut.
So that pulsar-perf produce / consume / read / transaction work transparently against both regular and scalable topics — V5 detects the topic kind via topic:// vs persistent:// lookup and routes accordingly. All four perf commands swap the v4 client API for the V5 surface: - PerformanceProducer: V5 Producer / AsyncProducer.async().send(), BatchingPolicy / ChunkingPolicy / CompressionPolicy / ProducerAccessMode, ProducerEncryptionPolicy backed by PemFileKeyProvider for --encryption-key-name + -value-file, V5 client-level TransactionPolicy + Transaction.async().commit(). - PerformanceConsumer: V5 QueueConsumer (all subscription types fold into QueueConsumer; Exclusive / Failover trigger a runtime warning). V5 has no MessageListener so each consumer gets a dedicated poll thread driving receive(Duration), invoking the same handler the v4 listener did. - PerformanceReader → CheckpointConsumer. Same poll-thread pattern as PerformanceConsumer. --start-message-id only accepts earliest / latest (V5 Checkpoint has no lid:eid form). - PerformanceTransaction: combined producer+consumer over V5 Transaction. Per-message ack is synchronous (V5 acknowledge is sync void). PerfClientUtils gets a createV5ClientBuilderFromArguments() factory mirroring the v4 helper. The TLS predicate only enables TLS when the user genuinely wants it (URL is pulsar+ssl, a trust cert path is supplied, or a flag is explicitly TRUE) — never on Boolean.FALSE from default-value resolution. The unconditional non-null check caused TLS to flip on against plaintext brokers. V5 features without a direct v4 equivalent are logged as no-ops or documented warnings: --max-outstanding(-across-partitions), --stats-interval-seconds, --max-lookup-request, --ssl-factory-plugin*, --busy-wait, --auto-scaled-receiver-queue-size, --batch-index-ack, --pool-messages, chunked-message knobs, --receiver-queue-size-across-partitions, --replicated, --use-tls. Per-transaction --txn-timeout moves to client-level TransactionPolicy.
…erBuilderImpl) CI build failed in PerformanceProducerTest.java:179 because the test cast createProducerBuilder()'s return to v4 ProducerBuilderImpl and inspected getConf().isBatchingEnabled(). The V5 ProducerBuilder is opaque (no public conf accessor), so that white-box assertion cannot survive the V5 migration. Regression intent — "disableBatching=true must propagate to the configured builder" — is now covered by V5's BatchingPolicy.ofDisabled() tests and the end-to-end perf workflow tests in this file.
dao-jun
approved these changes
May 29, 2026
Two V5 client-side bugs that the pulsar-perf migration tests exposed: 1. PulsarClientBuilderV5#authentication(plugin, params) set the plugin class name + params strings on the conf but never instantiated the Authentication and called conf.setAuthentication(...). The v4 PulsarClientImpl reads the Authentication instance via conf.getAuthentication() at connect time, so the client connected with no credentials and the broker rejected the handshake (caught by Oauth2PerformanceTransactionTest). Fix: instantiate via AuthenticationFactory.create and attach to the conf — matches what v4 ClientBuilder does in the same overload. Bad plugin class names are wrapped as V5 PulsarClientException rather than leaking the v4 exception type. New unit tests cover both paths. 2. QueueConsumerBuilder was missing replicateSubscriptionState. The underlying v4 ConsumerConfigurationData supports it and StreamConsumerBuilder already exposes it; QueueConsumerBuilder was the odd one out. Same one-line wiring through to the v4 conf. Needed by PerformanceTransactionTest.testTxnPerf which asserts on subscription.isReplicated().
Three CI-uncovered fixes on the perf side:
1. CmdBase: enable picocli case-insensitive enum parsing. V5 enums
(SubscriptionInitialPosition.EARLIEST, ProducerAccessMode.SHARED,
...) are uppercase, while users have been passing the v4 spellings
("Earliest", "Shared") on the perf CLI for years. Picocli is
case-sensitive by default and would silently fail to parse the
v4 spellings, exiting the perf command before any logs were
emitted. Caught by PerformanceTransactionTest.testConsumeTxnMessage.
2. PerformanceTransaction: collect per-txn send futures and await
them before committing. V5 transaction-aware sends are queued
onto an internal dispatch chain, so the v4-side txn-coordinator
registration of the send can race the commit() if commit fires
before the chain drains. Symptom:
InvalidTxnStatusException: ... unexpected state : COMMITTING,
expect OPEN. This is also the semantically correct ordering
(commit only after sends land). Caught by
Oauth2PerformanceTransactionTest.testTransactionPerf.
3. Wire --replicated through to the V5 QueueConsumerBuilder now
that it exposes replicateSubscriptionState (preceding commit).
Drop the "ignored" warning. Needed by
PerformanceTransactionTest.testTxnPerf.
…roducer Addresses review feedback: PerformanceProducer committed the transaction purely on the local send counter (numMessageSend == numMessagesPerTransaction), but messageBuilder.send() is asynchronous and V5 dispatches it through the producer dispatch chain. The commit could race ahead of the txn-coordinator registration of the sends, which the broker rejects with InvalidTxnStatusException (COMMITTING, expect OPEN). Collect the in-flight transaction's send futures and await them before committing — the same fix already applied to PerformanceTransaction. The send chain swallows per-send failures, so the join never throws on a send error.
void-ptr974
approved these changes
Jun 2, 2026
void-ptr974
left a comment
Contributor
There was a problem hiding this comment.
Thanks for the update. This makes V5 more practical to try and benchmark, especially with scalable topics.
Looking forward to the broader V5 client work.
LGTM.
This was referenced Aug 18, 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
lhotari
added a commit
that referenced
this pull request
Sep 8, 2026
…ient (#26480) ### Motivation `pulsar-perf` was migrated to the V5 client API in #25887 and `pulsar-client` in #25917, so every subcommand of both tools now drives the V5 SDK. Pulsar 5.0 does not deprecate ordinary (non-scalable) topics and the v4 client remains fully supported, so it is still useful to drive the CLI tools with the v4 client: for benchmark results comparable across broker versions (the v4 commands can also target pre-5.0 brokers, which the V5 SDK cannot talk to at all), and for testing the v4 client and ordinary topics without the V5 SDK in the path. The migration also left a set of options that are still accepted but no longer do what they say, with no way to get the behaviour back — round-robin partition routing and the outstanding-message limits on `produce`, the real subscription types and the receiver-queue knobs on `consume`, the v4 `Reader` and its `<ledgerId>:<entryId>` start position on `read`, the v4 transaction coordinator on `transaction`, KeyValue schemas and non-durable subscriptions on `pulsar-client`, and `loadConf` for `client.conf` keys without a dedicated CLI flag. ### Modifications Each command is split into an abstract base holding everything that is not client-specific, plus two thin sibling 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) AbstractCmdProduce -> CmdProduce / CmdProduceV4 (produce-v4) AbstractCmdConsumeCommand -> CmdConsume / CmdConsumeV4 (consume-v4) AbstractCmdReadCommand -> CmdRead / CmdReadV4 (read-v4) ``` `PulsarClientTool` builds a second, lazily-supplied v4 `ClientBuilder` for the `*-v4` commands; it keeps `loadConf` and accepts `http://` service URLs. `pulsar-shell` picks the new `pulsar-client` commands up for free, and both tools' `gen-doc` / `generate_documentation` render them. New subcommands only. The V5 commands keep their behaviour apart from a few deliberate fixes that the split surfaced (the `ProducerSocket` send/ack ordering and its `close()` null guard, the transaction receive-failure path), plus some added "use the `-v4` command" pointers on existing warnings — all listed in the PR description. ### Verifying this change Covered by `PulsarPerfTestToolTest`, `PerformanceV4CommandsTest`, `PerformanceTransactionV4Test`, `CmdV4CommandsTest`, `PulsarClientToolTest` (including the three KeyValue cases disabled by #25917, now re-enabled against `produce-v4`), and new `produce-v4` / `consume-v4` / `read-v4` cases in the `PerfToolTest` integration test.
Radiancebobo
pushed a commit
to Radiancebobo/pulsar
that referenced
this pull request
Oct 8, 2026
…ient (apache#26480) ### Motivation `pulsar-perf` was migrated to the V5 client API in apache#25887 and `pulsar-client` in apache#25917, so every subcommand of both tools now drives the V5 SDK. Pulsar 5.0 does not deprecate ordinary (non-scalable) topics and the v4 client remains fully supported, so it is still useful to drive the CLI tools with the v4 client: for benchmark results comparable across broker versions (the v4 commands can also target pre-5.0 brokers, which the V5 SDK cannot talk to at all), and for testing the v4 client and ordinary topics without the V5 SDK in the path. The migration also left a set of options that are still accepted but no longer do what they say, with no way to get the behaviour back — round-robin partition routing and the outstanding-message limits on `produce`, the real subscription types and the receiver-queue knobs on `consume`, the v4 `Reader` and its `<ledgerId>:<entryId>` start position on `read`, the v4 transaction coordinator on `transaction`, KeyValue schemas and non-durable subscriptions on `pulsar-client`, and `loadConf` for `client.conf` keys without a dedicated CLI flag. ### Modifications Each command is split into an abstract base holding everything that is not client-specific, plus two thin sibling 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) AbstractCmdProduce -> CmdProduce / CmdProduceV4 (produce-v4) AbstractCmdConsumeCommand -> CmdConsume / CmdConsumeV4 (consume-v4) AbstractCmdReadCommand -> CmdRead / CmdReadV4 (read-v4) ``` `PulsarClientTool` builds a second, lazily-supplied v4 `ClientBuilder` for the `*-v4` commands; it keeps `loadConf` and accepts `http://` service URLs. `pulsar-shell` picks the new `pulsar-client` commands up for free, and both tools' `gen-doc` / `generate_documentation` render them. New subcommands only. The V5 commands keep their behaviour apart from a few deliberate fixes that the split surfaced (the `ProducerSocket` send/ack ordering and its `close()` null guard, the transaction receive-failure path), plus some added "use the `-v4` command" pointers on existing warnings — all listed in the PR description. ### Verifying this change Covered by `PulsarPerfTestToolTest`, `PerformanceV4CommandsTest`, `PerformanceTransactionV4Test`, `CmdV4CommandsTest`, `PulsarClientToolTest` (including the three KeyValue cases disabled by apache#25917, now re-enabled against `produce-v4`), and new `produce-v4` / `consume-v4` / `read-v4` cases in the `PerfToolTest` integration test.
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
PIP-475 added regular-to-scalable topic migration on the broker and the V5 client SDK. The V5 client transparently routes against both regular and scalable topics. Pulsar's perf CLI (
pulsar-perf produce / consume / read / transaction) still used the v4 client API and therefore could not benchmark scalable topics at all.This PR migrates all four
pulsar-perfcommands to the V5 client API so the same binary now works against regular and scalable topics out of the box.A small V5 client-side fix is bundled in because it is required for the perf migration:
PulsarClientBuilderV5#tlsPolicywas a stub that dropped every TlsPolicy field on the floor (only settinguseTls=true). The perf commands take--trust-cert-fileand the related flags and expect them to land, so the policy is now wired through toClientConfigurationDataend-to-end.Modifications
pulsar-client-v5PulsarClientBuilderV5#tlsPolicynow wirestrustCertsFilePath,keyFilePath,certificateFilePath,allowInsecureConnection,enableHostnameVerificationthrough to the underlying v4 conf. Two new unit tests inPulsarClientBuilderV5Testcover field propagation and theTlsPolicy.ofInsecure()shortcut.pulsar-testclient— depends onpulsar-client-v5/pulsar-client-api-v5.PerfClientUtilsgetscreateV5ClientBuilderFromArguments()mirroring the v4 helper. Important detail: the TLS predicate only enables TLS when the user genuinely wants it (URL ispulsar+ssl://, a trust cert path is supplied, or a flag is explicitly TRUE). Earlier wiring that treated non-null Booleans as intent caused TLS to flip on against plaintext brokers because picocli / default-value resolution coerces those flags toBoolean.FALSEeven when not passed.PerformanceProducer: V5Producer/AsyncProducer.async().send().BatchingPolicy/ChunkingPolicy/CompressionPolicy/ProducerAccessMode.ProducerEncryptionPolicybacked byPemFileKeyProviderfor--encryption-key-name+--encryption-key-value-file. V5 client-levelTransactionPolicy+Transaction.async().commit().PerformanceConsumer: V5QueueConsumer. All subscription types fold intoQueueConsumersemantics (work distribution);Exclusive/Failovertrigger a runtime warning. V5 has noMessageListener, so each consumer gets one dedicated poll thread drivingreceive(Duration)and invoking the same handler the v4 listener did. V5acknowledge(...)is synchronous void — wrapped in try / catch with the existing failure counters.PerformanceReader→ V5CheckpointConsumer. Same poll-thread emulation asPerformanceConsumerreplacesReaderListener.--start-message-idnow only acceptsearliest/latest— V5Checkpointhas nolid:eidfactory and rejecting that form explicitly is cleaner than silently falling back toearliest.PerformanceTransaction: combined producer+consumer over V5Transaction. V5 puts.transaction(txn)on the message builder rather thanproducer.newMessage(txn). Same semantics, different ergonomics.V5 features without a direct v4 equivalent
The following CLI flags survive for back-compat but are logged as no-ops at start-up:
--max-outstanding,--max-outstanding-across-partitions,--stats-interval-seconds,--max-lookup-request,--ssl-factory-plugin*,--busy-wait,--auto-scaled-receiver-queue-size,--batch-index-ack,--pool-messages, chunked-message knobs (--max_chunked_msg,--expire_time_incomplete_chunked_messages,--auto_ack_chunk_q_full),--receiver-queue-size-across-partitions,--replicated,--use-tls.Per-transaction
--txn-timeoutmoves to client-levelTransactionPolicy.Verifying this change
PulsarClientBuilderV5Testcover the TlsPolicy wiring.This change is already covered by existing tests, so it does not need additional tests.
Does this pull request potentially affect one of the following parts:
pulsar-perfsubcommands. CLI flag UX is preserved end-to-end; behavior for flags listed above is logged as a no-op.