Skip to content

[improve][cli] Add v4-client subcommands to pulsar-perf and pulsar-client - #26480

Merged
lhotari merged 5 commits into
apache:masterfrom
lhotari:lh-improve-v4-cli-commands
Sep 8, 2026
Merged

lhotari merged 5 commits into
apache:masterfrom
lhotari:lh-improve-v4-cli-commands

Conversation

@lhotari

@lhotari lhotari commented Sep 7, 2026 •

Copy link
Copy Markdown
Member

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 be able
to drive the CLI tools with the v4 client:

  • Comparable benchmark results across broker versions. A large part of the value of
    pulsar-perf is comparing numbers between Pulsar releases. Those comparisons only hold if the
    client under test is the same one, so the v4 commands can also be pointed at pre-5.0 brokers,
    which the V5 SDK cannot talk to at all.
  • Testing the v4 client and ordinary topics without the V5 SDK in the path, so a V5-side change
    cannot influence the result.

Beyond that, the migration left a set of options that are still accepted but no longer do what they
say, with no way to get the behaviour back:

pulsar-perf

  • produce: --max-outstanding / --max-outstanding-across-partitions have no effect, and
    round-robin partition routing is gone.
  • consume: Exclusive / Failover / Key_Shared all degrade to work-queue semantics;
    --auto-scaled-receiver-queue-size, --batch-index-ack, --pool-messages,
    --receiver-queue-size-across-partitions and the chunked-message knobs are inert. Dispatch moved
    from the client's listener threads to per-consumer poll threads.
  • read: measures a V5 CheckpointConsumer, a different broker-side entity, so nothing exercises
    the v4 Reader. The <ledgerId>:<entryId> start position is rejected, and --use-tls and
    --receiver-queue-size are inert.
  • transaction: the v4 transaction coordinator stays a live broker code path next to the v5 one
    (PIP-473 P5.4) with nothing driving it. Acknowledgement latency is no longer measurable, because
    V5's acknowledge is a synchronous void, so the reported figure times a local call rather than the
    broker round trip.

pulsar-client

  • produce: --key-value-encoding-type is rejected outright and -kvk / -kvkf / -ks became
    dead flags; --disable-replication only warns.
  • consume: --subscription-mode NonDurable silently creates a durable subscription; --regex is
    not a regex any more but a namespace subscription over the pattern's tenant/namespace;
    --start-timestamp was dropped; -mc / -ac / -pm only warn.
  • read: the <ledgerId>:<entryId> start position is rejected (and its WebSocket base64 encoding is
    gone); -i and -q only warn.
  • Root: loadConf is gone, so every client.conf key without a dedicated CLI flag is silently
    ignored, and http:// / https:// service URLs no longer work.

Two PulsarClientToolTest cases were disabled by #25917 with a "deferred to a follow-up" note. This
is that follow-up.

Modifications

Rather than reverting or duplicating the commands, each one is split into an abstract base holding
everything that is not client-specific, plus two thin subclasses that bind the client types — the
existing V5 command and a new v4 one. The V5 and v4 classes are siblings, not parent and child.

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)

The bases hold the CLI options, the argument validation, the accounting and reports (perf) and the
WebSocket paths (client), so the improvements made since the migration are shared rather than
duplicated — instance recorders, the 3-significant-digit histograms and the latency clamps from
#26466, and the class-based subcommand registration from #26467. Each perf 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 rather than papered over:

  • 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 perf thread waits out the coordinator's asynchronous connect,
    so the rollover loop still counts every failed open;
  • pulsar-client's message rendering is not forced into one shape: the v4 output keeps the
    encryption context, the formatted broker publish and event times, the ordering key, the schema
    version and the index, none of which the V5 Message exposes.

PulsarClientTool builds a second, v4 ClientBuilder for the *-v4 commands. It keeps loadConf,
so every client.conf key still applies without a hand-written translation, and it accepts
http:// service URLs. It is supplied lazily, because updateConfig() is the preRun() hook and
runs for every invocation — building it eagerly would make --help, generate_documentation and
the V5 commands depend on a service URL and on client.conf keys only the v4 client parses.

pulsar-shell picks the new pulsar-client commands up for free, and both tools'
gen-doc / generate_documentation render them.

Also wires --jsse-provider / --jca-provider into the v4 client builder in PerfClientUtils,
which previously read them only for the admin and V5 legs.

The V5 commands are unchanged in behaviour except for the following, all of which are deliberate and
were verified line by line against the pre-split code:

  • ProducerSocket.send() now installs the completion future before sendText(...). The old order
    raced: an ack delivered on the Jetty read thread in between completed the previous future and left
    the new one uncompleted, so the caller blocked the full 30s and the produce loop aborted with -1.
    Latent since [pulsar-client-tools] Add support for websocket produce/consume command #3835.
  • ProducerSocket.close() gained a session != null guard. onClose nulls the session, so a
    remote-initiated close before close() used to NPE and turn a successful publish into exit -1.
  • The transaction receive-failure path catches Exception rather than PulsarClientException and
    returns after exit(1). The widening is forced — the v4 and V5 overrides throw two different
    PulsarClientException classes — and the return stops a fall-through into message.id() on a
    null message under a test exit procedure.
  • New user-visible output on nine sites: --queue-size is now reported as inert on read instead of
    silently ignored; the existing "no effect" warnings on produce / consume / read (both tools)
    and on transaction gained a "use the -v4 command" pointer; and pulsar-perf consume logs one
    new unconditional INFO line naming the V5 consumer API.

Verifying this change

This change added tests and can be verified as follows:

pulsar-testclient

  • PulsarPerfTestToolTest — subcommand registration and per-command naming, conf-file defaults and
    case-insensitive enums on the v4 commands, the shared -st enum mapping onto the v4 client enum,
    the logger naming the report lines depend on, and that a --delay-range draw of 0 is
    distinguishable from "no delay flag".
  • PerformanceV4CommandsTest — end-to-end produce / consume / read round trips through the v4
    commands against a mocked broker, partitioned-topic creation, the v4-only producer knobs asserted
    on the built configuration, the <ledgerId>:<entryId> start position (asserted positionally, so
    a reader that ignored it fails), and that read rejects that form while read-v4 accepts it.
  • PerfToolTest (integration, CLI group) gains produce-v4 / consume-v4 / read-v4 cases that
    run the commands from bin/pulsar-perf in a container.

pulsar-client-tools

  • CmdV4CommandsTest — subcommand registration and naming, case-insensitive enums on the v4
    subcommands, KeyValue schema building, the lid:eid start position including its WebSocket
    encoding, the --start-timestamp validation, and that commands needing no v4 client (--help,
    generate_documentation) still run without a service URL and with client.conf keys only the v4
    loader parses.
  • PulsarClientToolTest — the three KeyValue cases disabled by [improve][cli] Migrate pulsar-client to the V5 client API #25917 are re-enabled against
    produce-v4, and a new testNonDurableSubscribeWithV4Client asserts that consume-v4 really
    creates a non-durable subscription, i.e. it disappears when the consumer disconnects.

The full pipeline is green in Personal CI for both halves (43/43 checks each).

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

New subcommands only; no existing command's options or behaviour change.
pulsar-perf gains produce-v4, consume-v4, read-v4, transaction-v4, and pulsar-client
gains produce-v4, consume-v4, read-v4.

Documentation

  • doc
  • doc-required
  • doc-not-needed
  • doc-complete

The generated CLI reference on the website is produced from the commands themselves
(pulsar-perf gen-doc / pulsar-client generate_documentation), so the new subcommands are picked
up automatically.

…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.
… V5 ones

### Motivation

`pulsar-client` was migrated to the V5 client API in apache#25917, which left a set of
options that are still accepted but no longer do what they say:

- `produce`: `--key-value-encoding-type` is rejected outright and `-kvk` /
  `-kvkf` / `-ks` became dead flags; `--disable-replication` only warns.
- `consume`: `--subscription-mode NonDurable` silently creates a durable
  subscription; `--regex` is not a regex any more but a namespace subscription
  over the pattern's `tenant/namespace`; `Exclusive` / `Failover` degrade to
  work-queue semantics; `--start-timestamp` was dropped; `-mc` / `-ac` / `-pm`
  only warn.
- `read`: the `<ledgerId>:<entryId>` start position is rejected (and its
  WebSocket base64 encoding is gone), `-i` and `-q` only warn. `read` also
  drives a V5 `CheckpointConsumer`, a different broker-side entity, so nothing
  exercises the v4 `Reader`.
- Root: `loadConf` is gone, so every `client.conf` key without a dedicated CLI
  flag is silently ignored, and `http://` / `https://` service URLs no longer
  work.

Two `PulsarClientToolTest` cases were disabled by that migration with a
"deferred to a follow-up" note. This is that follow-up.

### Modifications

The same split as the pulsar-perf change: an abstract base per command holds
everything that is not client-specific — the CLI options, the argument
validation and the WebSocket path, which speaks HTTP and has no client
generation of its own — and two thin subclasses bind the client types.

    AbstractCmdProduce       -> CmdProduce  / CmdProduceV4  (produce-v4)
    AbstractCmdConsumeCommand-> CmdConsume  / CmdConsumeV4  (consume-v4)
    AbstractCmdReadCommand   -> CmdRead     / CmdReadV4     (read-v4)

`AbstractCmdConsume` keeps only what both generations share; the message
rendering genuinely differs (the v4 one prints the encryption context, the
formatted broker publish/event times, the ordering key, the schema version and
the index, none of which the V5 `Message` exposes), so it lives in
`V5MessageSupport` / `V4MessageSupport` rather than being forced into one shape.

`PulsarClientTool` builds a second, v4 `ClientBuilder` for the `*-v4` commands.
That one keeps `loadConf`, so every `client.conf` key still applies without a
hand-written translation, and it accepts `http://` service URLs. `pulsar-shell`
picks the new commands up for free, and so does `generate_documentation`.

### Verifying this change

New unit tests cover subcommand registration and naming, case-insensitive enums
on the v4 subcommands, KeyValue schema building, the `lid:eid` start position
(including its WebSocket encoding) and the `--start-timestamp` validation. The
three KeyValue `PulsarClientToolTest` cases disabled by apache#25917 are re-enabled
against `produce-v4`, and a new `testNonDurableSubscribeWithV4Client` asserts
that `consume-v4` really creates a non-durable subscription — the subscription
disappears when the consumer disconnects, which is exactly what the V5-based
`consume` cannot do.
@lhotari

lhotari commented Sep 7, 2026

Copy link
Copy Markdown
Member Author

This change is necessary to be able to profile the V4 client changes and ordinary V4 topic behavior in master branch.
The profiling test https://github.com/apache/pulsar/blob/master/tests/integration/src/test/java/org/apache/pulsar/tests/integration/profiling/PulsarProfilingTest.java would be updated so that both V4 and V5 pulsar-perf could be used in profiling.

@david-streamlio david-streamlio left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed the full diff (34 files, +6439/-3020) by theme rather than file order, focusing on the
central claim that "the V5 commands are unchanged in behaviour". I rebuilt the pre-split files and
compared them line by line against base + V5 subclass for all seven splits.

What I ran locally (JDK 25, macOS): quickCheck — pass; :pulsar-testclient:compileTestJava
and :pulsar-client-tools:compileTestJava — pass; PulsarPerfTestToolTest 6/6,
CmdV4CommandsTest 11/11, PerformanceV4CommandsTest 9/9. I did not run
PulsarClientToolTest or the PerfToolTest container ITs.

The split itself is the right shape — abstract base + two sibling subclasses, with the seams placed
where the clients genuinely differ, is much better than reverting or duplicating. Six of the seven
splits are faithful. The findings below are concentrated in one file.

Blocking

1. ProducerSocket.onMessage lost its null guard — AbstractCmdProduce.java:441-442

Base (CmdProduce.java:522-527):

log.info().attr("ack", msg).log("Received ack");
if (this.result != null) {
    this.result.complete(null);
}

Now:

log.info().attr("msg", msg).log("ack= ");
this.result.complete(null);

result is null until the first send(), so any text frame the WebSocket proxy sends outside a
pending send — before the first message, or after close() — now NPEs inside the Jetty
@OnWebSocketMessage callback and fails the session, where it was previously ignored. The
reordering in send() (finding 2) narrows the window but does not close it, because onMessage
fires for any text frame, not only acks.

This whole nested class looks copied from 820300f93b^ (pre-#25917) rather than moved from the PR
base — which would also explain the changed log line: every ack on the
pulsar-client produce ws://… path now logs ack= with attr msg instead of Received ack
with attr ack. Suggest restoring both the guard and the base's log text.

V5 behaviour changes that the description says do not exist

None of these look wrong — several are clear fixes — but the PR states the V5 commands are
unchanged, and a reviewer diffing later will trip over them. Worth a line in Modifications.

  1. AbstractCmdProduce.java:408-409 — send() now assigns this.result before sendText(...);
    the base assigned it after. The base ordering races: an ack arriving in between completes a
    stale future and the caller blocks the full 30 s. This is a real fix — keep it, disclose it.
  2. AbstractCmdProduce.java:450-451 — close() gains a session != null guard. Also a fix.
  3. PerformanceTransactionBase.java:457-469 — the receive-failure path now catches Exception
    (not just PulsarClientException) and returns after exit(1). Under a test exit procedure the
    base fell through into message.id() on a null message and NPE'd, so this is a fix too.
  4. PerformanceReaderBase.java:158-161 — FutureUtil.waitForAll(futures).get() is dropped in
    favour of draining future.get() in list order. Reader creation no longer fails fast on
    whichever future fails first. Minor, but it is a real difference.
  5. PerformanceConsumerBase.java:348-351 — the per-topic log line loses its consumerType
    attribute (moved to a one-time line). Changes what a per-topic grep sees.
  6. CmdRead.java:112-114 — new log when -q > 0; the base ignored receiverQueueSize silently.
    Together with the new "Use consume-v4 / read-v4 / produce-v4" sentences in the V5 warnings
    (CmdProduce.java:70, CmdConsume.java:85,94, CmdRead.java:88-89), this is new user-visible
    output on commands described as unchanged. All sensible additions — just not "unchanged".

Test coverage

8. transaction-v4 has no functional test at all. PerformanceTransactionV4 is 144 lines of
new production code, and every test reference is metadata-only: registered under its own name
(PulsarPerfTestToolTest.java:51), its flag set is a subset of transaction
(:79-81), and its logger is named correctly (:115-116). Nothing opens, commits or aborts a
transaction through it. I grepped every test source tree to confirm. Driving the v4 transaction
coordinator is the stated motivation for the command, and the V5 counterpart has three real tests
in PerformanceTransactionTest. This is the one new subcommand of seven with zero behavioural
coverage.

9. consume-v4's --start-timestamp / --end-timestamp is only tested negatively.
CmdConsumeV4.java:134-135 calls consumer.seek(startTimestamp) and :147 filters on
getPublishTime(). The only coverage is two rejection messages
(CmdV4CommandsTest.java:218-229). Delete the seek(...) and the filter and every test in this PR
still passes — yet CmdV4CommandsTest.java:69 names timestamp seek as the reason consume-v4
exists.

Test quality

  1. PerformanceV4CommandsTest.java:169-177 — asserts getDeliverAtTime() positionally on two
    received messages, assuming the -dr 0,1 message arrives first. A message stamped
    deliverAtTime = now goes through the delayed-delivery tracker; one tick of deferral swaps the
    order and the test fails for an unrelated reason. The payloads are already distinct
    ("delayed" / "plain") — key the assertions off new String(msg.getData()) instead.
    (Passed 9/9 locally, but the race is structural.)
  2. PerformanceV4CommandsTest.java:213 — asserts backlog with no Awaitility. exitLatch is
    released from PerfClientUtils.exit(0) (PerformanceConsumerBase.java:483) before
    stopConsuming(); closeClient(client) flushes pending acks, and acks are grouped on a 100 ms
    delay. Awaitility.untilAsserted is what CODING.md asks for here.
  3. PulsarClientToolTest.java:578,629 — the two re-enabled cases carry no timeOut (every other
    @Test in the file does), and the CompletableFuture they create is completed exceptionally at
    :648 but never inspected, so a non-zero produce-v4 exit surfaces as assertNotNull(message)
    with the cause discarded. testNonDurableSubscribeWithV4Client:216-217 gets this right.
    To be clear, the three re-enabled KeyValue cases are otherwise genuine — they decode with a
    Schema.KeyValue consumer and would fail against a broken produce-v4.
  4. TestCmdConsume.java:36-38 — the PR already retargets this reflection to
    AbstractCmdConsumeCommand. Since subscriptionName is protected
    (AbstractCmdConsumeCommand.java:83) and the test is in the same package, the three lines and
    the java.lang.reflect.Field import can just be deleted in favour of
    cmdConsume.subscriptionName = "my-sub"; — CODING.md's no-reflection rule, and this is the
    moment since the line is being touched anyway.

v4-only correctness

  1. PerformanceTransactionBase.java:522-532 + PerformanceTransactionV4.java:119-124 — on V5 the
    ack callback runs inline on the worker thread, but v4's acknowledgeAsync completes on a client
    internal thread. The exceptionally block interrupts that thread and takes the early
    return null, skipping numMessagesAckFailed.increment(). So on transaction-v4 an
    interrupt-caused ack failure is silently uncounted and the worker never learns to stop. The same
    pattern in sendAndRecord is pre-existing; the ack side is new.
  2. CmdProduceV4.java:189-194 — nativeAvroSchemaOrNull returns null when getNativeSchema() is
    empty, and generateMessageBodies then ships the raw JSON text as the payload. The historic v4
    code did nativeSchema.get() and threw. AUTO_PRODUCE_BYTES delegates getNativeSchema(), so
    the normal -s "avro:{…}" path is unaffected — but the fallback turns a loud failure into
    silently wrong data. Prefer throwing.

Nits

  1. AbstractCmdReadCommand.java:96 — protected abstract String startMessageId() is implemented
    twice and never called (webSocketStartMessageId() is the one used at :140). Dead; delete.
  2. PulsarClientTool.buildV4ClientBuilder — the "proxy-protocol must be provided with proxy-url"
    throw is unreachable; updateConfig() already returns 1 for that case in preRun().
  3. PulsarClientTool.buildV4ClientBuilder — serviceUrl(...) and tlsTrustCertsFilePath(...) are
    applied unconditionally after loadConf(conf), so a null --tlsTrustCertsFilePath nulls out
    the value client.conf just supplied. Faithful to the pre-#25917 code, so not a new
    regression — but since it is being written fresh, isNotBlank guards would fix a latent bug.
  4. PerfClientUtils.java:134-148 — the new v4 block mirrors the admin block exactly, good. Note
    the V5 path applies the providers only inside if (wantsTls(arguments)) while v4/admin apply
    them unconditionally, so --jsse-provider on a plaintext URL is a silent no-op on V5 and a
    silent write on v4. Harmless, but inconsistent for a flag whose purpose is FIPS validation.
  5. CmdV4CommandsTest.java:231-234 uses fully-qualified java.util.Set / TreeSet / Arrays
    inline; the sibling PulsarPerfTestToolTest.java:176-181 imports them.

Verified clean

For the record, these all checked out: every LongAdder in the producer, consumer and transaction
splits increments on exactly the same paths (the commit/abort duplication collapsing into
endTransaction accounts for the whole delta); nextDeliverAfterSeconds() preserves the original
delay > 0 / delayRange precedence, and returning Long rather than a sentinel is the right call
for a drawn delay of 0; V5MessageSupport is character-identical to the base rendering code, so no
printed output changes on the V5 path; the lazily-supplied v4 ClientBuilder is reached only by the
three *-v4 commands and each wraps it in try-with-resources, so --help and
generate_documentation still run without a service URL and the client is never leaked; the
@CustomLog → Logger.get(getClass()) change is name-preserving, so PerfToolTest's aggregated-
throughput assertions still match; messageFormatter and executorShutdownNow going static →
instance is a genuine improvement; all new files carry the ASF header and quickCheck is clean.

@david-streamlio

Copy link
Copy Markdown
Contributor

One follow-up from reading #26466, which this PR builds on. Not a defect in either change — the rationale comments for LATENCY_HISTOGRAM_SIGNIFICANT_DIGITS describe the world as it was before #26466, and this PR carries the wording forward unchanged:

  • PerfClientUtils.java:61 — "...and since these recorders are static fields the JVM pays for every perf subcommand, not just the one being run."
  • PerfClientUtilsTest.java:237 — "The perf clients hold their latency recorders in static fields, so every subcommand allocates them on startup rather than only the one being run."

#26466 converted every one of those to an instance field in the same commit that added these comments — git grep -E "static.*Recorder [a-zA-Z]" -- pulsar-testclient/src/main returns nothing at bc1785dcd8, and it still returns nothing here. So the premise the paragraphs rest on is the exact thing that commit fixed.

The conclusion is untouched — 3 digits is still clearly right, and the ceiling the test pins is still worth pinning. It is only the "every subcommand pays" justification that no longer applies; the cost is now bounded to the one command being run. Measured on HdrHistogram 2.1.9, for the ranges the test uses:

range 5 digits 3 digits
1 h in µs 16.00 MB 0.18 MB
10 d in ms 14.00 MB 0.16 MB

(The "11-22 MB" figure holds if you include the pre-#26466 33 h range, which came to 21.00 MB.)

Since PerfClientUtils.java is already touched here, it may be easiest to reword both in this PR — something like "a single Recorder over the ranges used here costs 14-16 MB at 5 digits, per running command" — rather than leaving a comment that contradicts the field declarations a few lines away. Equally fine as a separate cleanup if you would rather not widen this one.

…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.
…mmands

Follow-up to the review on apache#26480, pulsar-client half.

- `ProducerSocket.onMessage` had lost the `if (this.result != null)` guard, so a
  text frame arriving outside a pending send would NPE on the Jetty callback
  thread where it used to be ignored. Restores the guard and the pre-split log
  text ("Received ack" with attr `ack`); the version in the split was a
  hand-rendering of the pre-apache#25511 slf4j format string. Also drops the dead
  `throws InterruptedException` on `onConnect`, and comments the two deliberate
  changes in the same class: `send()` installs the completion future before
  `sendText(...)` (the old order let an ack complete the previous future and
  hang the caller for the full 30s), and `close()` null-checks the session that
  `onClose` clears.

- `consume-v4`'s `--start-timestamp` seek and `--end-timestamp` filter were only
  covered by their rejection messages: deleting either line left every test
  passing, though timestamp seek is the stated reason the command exists. Adds
  `PulsarClientToolTest.testConsumeV4StartAndEndTimestamp`, which pins both
  against a real observed publish time. Mutation checked: deleting either line
  fails it.

- `CmdProduceV4.nativeAvroSchemaOrNull` returned null for an Avro schema with no
  native definition, which would ship the raw JSON text as the payload; the
  pre-apache#25917 command threw. Throws again. Not reachable from any current CLI
  invocation, but the invariant is worth keeping explicit.

- `PulsarClientTool.buildV4ClientBuilder`: removes the unreachable
  "proxy-protocol must be provided with proxy-url" throw — `updateConfig()` is
  the `preRun()` hook and already returns 1 for that case before the supplier
  that reaches this method is even created, and `ClientBuilderImpl` re-checks it.
  Keeps the outer `isNotBlank` guard, which unlike the V5 path is load-bearing
  here because this builder uses `loadConf()`. Documents why `serviceUrl` and
  `tlsTrustCertsFilePath` are applied unconditionally.

- Deletes `AbstractCmdReadCommand.startMessageId()` and both overrides: never
  called, `webSocketStartMessageId()` is the one in use.

- `TestCmdConsume` sets `subscriptionName` directly instead of reflecting into
  it; the field is protected and the test is in the same package.

- `CmdV4CommandsTest` imports the `java.util` types its helper used inline.
@lhotari

lhotari commented Sep 7, 2026

Copy link
Copy Markdown
Member Author

Thank you for this — rebuilding the pre-split files and diffing them against base + V5 subclass for
all seven splits is exactly the check this PR needed, and it found a real regression I had missed.
I have addressed everything below. Three findings I could not reproduce and am pushing back on with
mechanism; a few others were right on the fact but wrong on the impact, and I have noted where,
because I think the corrected version is worth having on the record.

Blocking

1. ProducerSocket.onMessage lost its null guard — fixed.
You are right, and it was my mistake. Restored both the if (this.result != null) guard and the
base's log.info().attr("ack", msg).log("Received ack").

One correction on provenance: the class was not copied from 820300f93b^. That revision has both
the guard and "Received ack" — in fact every historical revision has the guard.
git log -S'ack= ' -- .../CmdProduce.java returns exactly two commits: #3835 (2019), which
introduced the slf4j LOG.info("ack= {}", msg), and #25511, which converted it to slog. So what I
wrote was a hand-reconstruction of the pre-#25511 format string with the guard dropped, not a
copy of any single revision. The onConnect(Session) throws InterruptedException in the same class
came from ConsumerSocket, which confirms it was reassembled by hand. I have dropped that dead
throws too.

Also a small correction on reachability, which is why I have not treated it as blocking for the
release: against Pulsar's own WebSocket endpoint the window is unreachable — every
sendAckResponse(...) call site in ProducerHandler is inside onWebSocketText or the sendAsync
continuation it starts, and error paths go through session.close(...), i.e. @OnWebSocketClose.
Only a non-Pulsar endpoint behind --url ws://… can send an unsolicited text frame. And "or after
close()" does not hold: by then send() has assigned result at least once, so a late frame
completes an already-completed future, which is a no-op rather than an NPE. It is still a real
fidelity regression in a commit presented as a move, so it is fixed regardless.

V5 behaviour changes not disclosed in the description (2, 3, 4, 6, 7)

All correct, and the description was wrong to claim "unchanged in behaviour" without qualification.
I have added a Behaviour changes on the V5 commands subsection to Modifications (text in Part 3
below) listing every one of them. Kept as they are — they are fixes — with these notes:

  • 2 (send() reordering) — kept. Worth spelling out that it and finding 1 are the same bug seen
    from two sides: the result != null guard was the band-aid over exactly this window, so the fix is
    "reorder and keep the guard", not one or the other. result is volatile, the caller is
    strictly serialised, and ProducerHandler sends exactly one ack per inbound message and nothing
    unsolicited, so the new ordering introduces no window of its own.
  • 3 (close() null guard) — kept, disclosed.
  • 4 (broadened catch) — kept, disclosed, with two corrections. The fall-through was not
    reachable in practice at base: both existing exit procedures
    (PerformanceTransactionTest:74-79, Oauth2PerformanceTransactionTest:112-117) throw on a
    non-zero code and production uses System::exit; and under a returning procedure the
    message.id() NPE landed inside the pre-existing catch (Exception ackEx), so it was miscounted
    as an ack failure rather than escaping. Also, widening to Exception is forced, not
    discretionary — the V4 and V5 overrides throw two different PulsarClientException classes
    (o.a.p.client.api vs o.a.p.client.api.v5), so no narrower shared catch exists, and the try body
    is only receive(consumer).
  • 6 (consumerType attribute) — restored. scalableConsumerType is V5-only so it cannot live in
    the shared base, but the base can ask for it: consumerTypeForLog() returns the shared
    subscriptionType by default and PerformanceConsumer overrides it with scalableConsumerType.
    The V5 line is byte-identical to base again, and consume-v4 gets the equivalent.
  • 7 (new user-visible output) — disclosed, and the list is longer than the four you found. The
    same sweep also picks up CmdProduce.java:80-81, PerformanceConsumer.java:94,
    PerformanceReader.java:63 and PerformanceTransaction.java:104, plus a brand-new
    unconditional INFO line at PerformanceConsumer.java:88 ("Using V5 scalable-topic consumer API")
    on every pulsar-perf consume run. Nine sites, all listed in the description now.

Findings I could not reproduce (5, 10, 18)

I checked each of these three carefully and I do not think the mechanism exists. Happy to be shown
wrong.

5. PerformanceReaderBase dropping FutureUtil.waitForAll — no behavioural difference.
FutureUtil.waitForAll is CompletableFuture.allOf (FutureUtil.java:56-61), and allOf is not
fail-fast: BiRelay.tryFire returns without firing unless both operand results are already present,
so the composed future settles only after every input has settled. The drain loop therefore reports a
failure no later than the removed waitForAll(...).get(), not later. The throwable is identical
too — allOf's BiRelay tests the left operand first and so reports the leftmost failure, which is
also where the in-order drain throws, and reportGet unwraps the CompletionException, so both
surface ExecutionException(cause). Neither version closes already-created readers on failure, so no
leak changed. And PerformanceConsumer.java:328-331 already used the bare drain loop before this PR
— PerformanceReader.java:139 was the module's only waitForAll caller — so this is a
normalisation onto the existing pattern rather than a dropped guard. I tried restoring it and then
reverted; happy to add it back if you still want the symmetry, but it would make failure reporting
strictly later.

10. Positional getDeliverAtTime() assertions — the ordering hazard is not reachable.
The consumer at :154 uses the default Exclusive subscription type
(ConsumerConfigurationData:90), which is served by PersistentDispatcherSingleActiveConsumer — and
trackDelayedDelivery is overridden only by PersistentDispatcherMultipleConsumers and its Classic
variant, so the single-active dispatcher inherits the no-op default at Dispatcher.java:118 and the
delayed-delivery tracker is never consulted. Even on a Shared subscription it would not defer:
InMemoryDelayedDeliveryTracker.addMessage rejects deliverAt <= getCutoffTime()
(clock.millis() + tickTimeMillis), and deliverAfter(0, SECONDS) is already due at send time —
"one tick of deferral" is precisely the window that cutoff excludes. With one non-batching producer
on a non-partitioned topic and each send(...).get() awaited, publish order is fixed and arrival
order follows it.

That said, your instinct that the test reads fragilely is fair, so I took the half of the suggestion
that costs nothing: the assertions now also check the payload identity of each message
(assertThat(new String(first.getData(), UTF_8)).isEqualTo("delayed")), with a comment recording why
the order is deterministic. A future reordering now fails as "expected delayed, got plain" instead of
as a confusing deliverAtTime mismatch, and the incidental ordering check is still pinned.

18. isNotBlank guards in buildV4ClientBuilder — the divergence cannot occur.
rootParams.tlsTrustCertsFilePath is not independent of client.conf; it is filled from it. The
option declares descriptionKey = "tlsTrustCertsFilePath" (PulsarClientTool.java:92-95) and the
commander's default-value provider is PulsarClientPropertiesProvider extends PropertiesDefaultProvider over the very same Properties object that becomes the loadConf map
(:138-139, :176, :228). So line 257 rewrites the value loadConf just applied, or applies the
genuine CLI override; it can only be null when the conf key is absent too, i.e. null over null.
Adding a guard would also remove the ability to clear a conf-supplied path with
--tlsTrustCertsFilePath "".

serviceUrl is the opposite of the suggested change: v4 ClientBuilderImpl.serviceUrl rejects a
blank value with checkArgument rather than nulling the field, and the call is load-bearing because
loadConf only reads the serviceUrl key while the shipped conf/client.conf supplies
brokerServiceUrl / webServiceUrl (which PulsarClientPropertiesProvider.create synthesises).
Guarding it would turn a clear "Param serviceUrl must not be blank" into a vaguer failure from
build(). I have left both unconditional and added a comment saying why, since it clearly reads as a
bug otherwise.

Test coverage (8, 9)

8. transaction-v4 had no functional test — added, and mutation-verified.
Correct, and this was the right thing to catch. New PerformanceTransactionV4Test (reusing
PerformanceTransactionTest's coordinator setup but on plain persistent:// topics, verified with
the v4 SDK client the base test provides):

  • testTransactionV4CommitsWhatItProduced — a committed run makes what it produced visible and
    drops the consume backlog by one per transaction, so the acks really were committed.
  • testTransactionV4AbortUndoesBothTheSendAndTheAck — -abort leaves nothing visible on the produce
    topic, leaves the consume backlog at 10, and leaves all ten messages immediately redeliverable on
    the tool's own subscription. That last part is what separates a real abort from a transaction that
    was merely never ended (whose acks would stay pending until the timeout).

Both pass. Mutation-tested to confirm they are load-bearing:

mutation result
sendMessage ignores the transaction killed
acknowledgeAsync ignores the transaction killed
acknowledgeAsync becomes a no-op (never acks at all) killed
abortTransaction becomes a no-op killed
drop .sendTimeout(0, TimeUnit.SECONDS) survived

I want to be straight about that last row rather than claim more than I verified — and chasing it
turned up a stale comment of my own. .sendTimeout(0, TimeUnit.SECONDS) was annotated in this PR
with "a send timeout and a transaction are mutually exclusive on the v4 client". That was true:
ProducerBase.newMessage(Transaction) threw "Only producers disabled sendTimeout are allowed to
produce transactional messages" — until bbf2a47867f "[fix][txn] Allow producer enable send timeout
in transaction (#16519)" deleted the check in 2022. Today the method only does
checkArgument(txn instanceof TransactionImpl).

So the line is still right, for a weaker reason: a send timeout is now permitted, but it still fails
the send on its own schedule and takes the transaction with it (that is exactly what #16519's own
testSendTxnMessageTimeout asserts), so disabling it leaves --txn-timeout as the only deadline.
Both stale comments (PerformanceTransactionV4 and PerformanceProducerV4) now say that instead.
The line stays unpinned by the test because the failure only appears once a send actually times out,
which five quick transactions never do — pinning it would need a deliberately stalled transaction,
which I did not think was worth the CI minutes.

One correction to the finding: the V5 counterpart has one real test for the transaction command
(PerformanceTransactionTest:139), not three — the other two drive PerformanceProducer and
PerformanceConsumer with -txn. Lower baseline, same conclusion.

9. consume-v4's timestamp flags only tested negatively — added, and mutation-verified.
Correct, and your mutation claim was right: both lines are dead under the defaults
(startTimestamp = 0L fails the > 0L guard, endTimestamp = Long.MAX_VALUE can never be
exceeded). New PulsarClientToolTest.testConsumeV4StartAndEndTimestamp publishes two messages, reads
their real publish time back as the boundary, waits for the millisecond clock to cross it, publishes
two more, then asserts -stp boundary+1 from -p Earliest prints only the newer pair and
-etp boundary with -n 4 prints only the older pair.

mutation result
delete consumer.seek(startTimestamp) killed
delete the getPublishTime() > endTimestamp break killed

Worth noting -etp has been equally untested since #24521 — this PR restores the -stp seek that
#25917 dropped, it did not introduce the gap.

Test quality (10, 11, 12, 13)

  • 10 — see above.
  • 11 (backlog assert without Awaitility) — wrapped in Awaitility.untilAsserted. Your stated
    mechanism is not quite the one that bites, though: the assert is not reached straight off the
    latch, because stop(thread) joins the perf thread and that thread cannot die until
    PerfClientUtils.closeClient has run a full client.close(), which flushes the grouped acks and
    blocks on the CloseConsumer round trip the broker can only answer after applying them. On the
    happy path it is ordered. What is genuinely unguarded is the unchecked 30 s join and the
    interrupt reportDone delivers after the latch, which can abort that close and delete the only
    happens-before edge. Awaitility closes both, and costs nothing on the ordered path — so the fix is
    right, for a different reason.
  • 12 (timeOut and the discarded cause) — both fixed: timeOut = 60000 on each, and future.get()
    before the receive, so a non-zero produce-v4 exit surfaces as its own exception instead of as
    assertNotNull(message) ten seconds later. One correction: "every other @Test in the file"
    carries a timeOut is not so — :74, :492 and :512 are bare @Test, and :512 blocks on a
    10 s receive just like these. And the PR did not introduce the swallowing; the bodies are
    unchanged from base apart from produce → produce-v4, and the same pattern is in the untouched
    testProducePartitioningKey. Fixed here anyway, because this PR is what makes these two execute
    in CI again.
  • 13 (reflection in TestCmdConsume) — done: the three lines, the java.lang.reflect.Field
    import and the now-unnecessary throws Exception on setUp are all gone, replaced by
    cmdConsume.subscriptionName = "my-sub";. Note the same pattern in TestCmdRead:24 is not
    fixable this way — startMessageId is still private in CmdRead.

v4-only correctness (14, 15)

  • 14 (interrupt on the wrong thread) — fixed by capturing the worker thread and interrupting
    that, so V5 is unchanged (there the captured thread is the current thread) and v4 no longer
    interrupts a client-owned executor thread. Your mechanism is exactly right on all three sub-claims.
    One correction on impact: I could not find any path that completes a v4 ack future with an
    InterruptedException in its cause chain — every such catch in ConsumerBase/ConsumerImpl is in
    the synchronous .get() wrappers — so it is a latent wrong-thread hazard rather than an observable
    miscount today. Still worth fixing since the ack-side callback is new here.

    On reflection I applied it to all three callbacks, not just the ack: sendAndRecord and
    endTransaction in the same file interrupt Thread.currentThread() from futures that complete on
    client-internal threads too, and all three are invoked from the same runWorker loop, so the
    capture is equally valid. You noted the send side is pre-existing and it is — but leaving two of
    three wrong next to a comment explaining why the third is right just invites the question. It is
    four lines and no behaviour change on V5.

  • 15 (nativeAvroSchemaOrNull silently shipping raw JSON) — changed to throw. One correction:
    it is not reachable from any CLI invocation today. The only AVRO-typed schema CmdProduceV4 can
    build is Schema.AUTO_PRODUCE_BYTES(Schema.generic(…AVRO)), and AvroBaseStructSchema.getNativeSchema()
    is a hard Optional.of(schema); a malformed avro: definition already throws in parseAvroSchema
    during construction. And in the counterfactual, AUTO_PRODUCE_BYTES validates on encode, so raw
    JSON would throw SchemaSerializationException at send() — a different loud failure, not corrupt
    data on the topic. So it is an invariant worth keeping explicit rather than a live defect; I used
    IllegalStateException to match that (the argument-level errors in this class use
    IllegalArgumentException).

Nits (16, 17, 19, 20)

  • 16 (dead startMessageId()) — deleted, all three sites. Confirmed no callers repo-wide; not
    picocli-bound and not used reflectively (TestCmdRead reflects on the field).
  • 17 (unreachable proxy-protocol throw) — deleted. Confirmed unreachable: updateConfig() is the
    preRun() hook, the execution strategy runs it before RunLast and returns early on a non-zero
    code, and the lazy supplier is only created after that guard. It is redundant a second time over
    — ClientBuilderImpl.proxyServiceUrl already does checkArgument(proxyProtocol != null, …). I
    kept the outer isNotBlank guard: unlike the V5 path this builder uses loadConf(), so passing
    nulls through would clear a proxyServiceUrl / proxyProtocol coming from client.conf.
  • 19 (--jsse-provider gated on V5, unconditional on v4/admin) — real observation, but I do not
    think either direction of alignment is an improvement, so I have left it. Removing the wantsTls
    gate is not possible: PulsarClientBuilderV5#tlsPolicy sets useTls=true for CLIENT_DEFAULT, so
    a provider flag on a pulsar:// URL would force a TLS handshake against a plaintext broker — the
    exact bug that gate and its javadoc exist to prevent. Gating the v4/admin writes instead is also
    wrong: wantsTls inspects arguments.serviceURL, not the separate adminUrl the admin builder is
    handed, and the v4 write is not always inert —
    ClientTlsFactorySupport.foldOAuth2IdpPolicy inherits conf.getJsseProvider() /
    getJcaProvider() into the CLIENT_OAUTH2 IdP policy, and needsClientTlsFactory() builds the
    factory on a plaintext broker when the OAuth2 plugin carries IdP TLS material. If anything it is
    the V5 side that drops a meaningful pin in that narrow case, and closing it would need
    tlsPolicy(CLIENT_OAUTH2, …), which would clobber the plugin's own IdP trust via putIfAbsent.
    That is a PIP-478 question, not this PR's.
  • 20 (fully-qualified java.util types) — fixed; imported Arrays, Set, TreeSet and dropped
    the qualifiers. A scan of all 34 changed files found no other occurrence.

Follow-up comment: stale LATENCY_HISTOGRAM_SIGNIFICANT_DIGITS rationale

Right on both counts, and thank you for catching that the premise was the very thing #26466 fixed.
Every Recorder under pulsar-testclient/src/main is an instance field at HEAD (15 declarations),
and the static→instance conversion happened in the same commit that added those comments.

I re-measured against HdrHistogram 2.2.2 and your figures reproduce exactly — 16.00 MB / 0.18 MB for
1 h in µs, 14.00 MB / 0.16 MB for 10 d in ms; the old "11-22 MB" is what the pre-#26466 ranges
(SECONDS.toMicros(120000) / toMillis(120000)) produce, which is the other half of the same
staleness. Both comments are reworded to "14-16 MB", with the per-command cost expressed as "a
command holds several of them (a live and a cumulative recorder per measured latency)" rather than
the old "every subcommand pays". The 3-digit conclusion and the 512 KiB test bound are unchanged.

Not addressed

Nothing. Every numbered finding plus the follow-up comment is either fixed, disclosed in the
description, or answered above with the mechanism.

…ir run budget

The bounded join in PerformanceTransactionV4Test.run() allows 90s, but the
method timeOut was 120s, leaving only 30s for assertions that themselves allow
up to 30s of Awaitility plus ten 15s receives. A slow run would trip TestNG's
opaque timeout instead of the helper's explicit "did not finish within"
failure. Raise the method budget to 180s so the helper reports first.
@lhotari
lhotari merged commit 31dea17 into apache:master Sep 8, 2026
81 of 83 checks passed
@lhotari lhotari added this to the 5.0.0-M2 milestone Sep 12, 2026
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.
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.

3 participants