Skip to content

[feat][fn] PIP-496: Support the V5 client and topic:// topics in Pulsar Functions and IO - #26697

Merged
lhotari merged 11 commits into
apache:masterfrom
lhotari:lh-pip-496-fn-v5-client-impl
Sep 30, 2026
Merged

lhotari merged 11 commits into
apache:masterfrom
lhotari:lh-pip-496-fn-v5-client-impl

Conversation

@lhotari

@lhotari lhotari commented Sep 23, 2026 •

Copy link
Copy Markdown
Member

PIP: PIP-496, proposed in #26698

Motivation

Scalable topics (topic://, PIP-460) can only be produced to and consumed from with the V5 client
(PIP-466). The Pulsar Functions and IO runtime uses the v4 client throughout, so today a function, source or
sink cannot read from or write to a topic:// topic.

It is also a migration blocker. PIP-475 refuses (HTTP 409) to migrate a persistent:// topic to a scalable
topic while any v4 producer or consumer is attached, and every function or connector on the topic attaches
one. A component needs a way to use the V5 client for persistent:// topics too, so that its topics can be
migrated.

This PR implements PIP-496.

Modifications

Client selection. Each component uses one client for its own topics. The runtime picks it with the same
rule as the pulsar-client and pulsar-perf CLIs (#26693):

  • A topic:// topic selects the V5 client, and any other domain selects v4.
  • An explicit clientApi=V5 makes the V5 client drive persistent:// topics.

The new clientApi setting (V4/V5; unset means "pick from the topics") is added to FunctionConfig,
SourceConfig and SinkConfig, as proto field FunctionDetails.clientApi = 24 (AUTO by default) and as
--client-api on pulsar-admin functions|sources|sinks. ClientApiResolver validates it when a component
is converted, which is when the worker registers or updates it, and again when an instance starts. It rejects:

  • mixing topic:// with other domains, or naming a segment:// topic directly;
  • V4 with a topic:// topic, and V5 with a non-persistent:// topic;
  • with V5: topic patterns, non-Java runtimes, and EFFECTIVELY_ONCE with an output topic (per-segment
    sequence IDs are not monotonic across splits and merges);
  • with V5 on persistent:// inputs: any subscription type other than Shared, because a V5 stream consumer
    cannot yet subscribe to a topic that has not been migrated;
  • with V5: consumer and producer encryption, message payload processors, consumer properties and
    skipToLatest;
  • with V5 and an ordered subscription: maxMessageRetries, deadLetterTopic and timeoutMs, which a stream
    cannot honor.

Instances resolve again when they start, so later versions may relax these rules but must not tighten them.

Consuming (V5PulsarSource). Each input topic gets its own V5 consumer.

  • A Shared subscription uses a QueueConsumer. Individual acks, negative acks, the processing timeout, the
    negative-ack delay and the dead letter policy map from the component's settings. Without a configured dead
    letter topic, the v4 default name (<input>-<subscription>-DLQ) is used, so that switching a component to V5
    keeps its dead letter topic.
  • Failover, Key_Shared and EFFECTIVELY_ONCE use a StreamConsumer, which delivers in order and splits key
    ranges across the component's instances.
    • A stream acknowledges only cumulatively, so StreamAckTracker acknowledges the completed prefix of the
      received records. Records completed out of order are never skipped.
    • A cumulative acknowledgment of a record completes every record received before it.
    • A stream has no negative acknowledgment, so a failed record fails the instance through its fatal handler,
      from whichever thread calls fail().
  • Each start of an instance uses a new consumer name, <tenant>-<namespace>-<name>-<instanceId>-<random>: a
    stream consumer group treats a reused name as a reconnect of the previous member.
  • Records are PulsarRecords over the v4 message that the V5 message wraps. Connectors that read the
    schema, the message or the encryption context keep working. The record's topic is the topic:// input,
    not the segment.

Publishing. V5 producers are presented through the v4 Producer/TypedMessageBuilder/MessageId
interfaces, so PulsarSink, ProducerCache, ContextImpl.newOutputMessage and the log appender keep a
single code path; only producer creation branches (V5ProducerFactory). Context.newOutputMessage publishes
to a topic:// topic with the V5 client in every component.

Runtime.

  • The thread runtime factory owns one V5 client. It is created on first use, shared by the instances, and
    closed with the factory.
  • InstanceUtils.createPulsarClientV5Builder applies the service URL, authentication, TLS and memory limit
    of the v4 client. The TLS policy is set only when TLS is enabled, and hostname verification is always
    passed explicitly, because the TLS policy's default differs from the v4 client's.
  • Without a configured memory limit, the V5 client keeps its 64 MiB default. The V5 producer has no
    pending-message limit, so the memory limit is what bounds it.
  • The V5 client implementation is looked up once under the runtime's class loader, so that the first V5 call
    from user code in the process and Kubernetes runtimes cannot fail the JVM-wide lookup.

User API (PIP-496): one default method, BaseContext.getPulsarClientV5(), which returns the runtime's
shared V5 client.

Packaging and worker.

  • pulsar-functions-api depends on pulsar-client-api-v5, so the V5 API is in java-instance.jar. The
    OpenTelemetry API is kept out of that jar, as with the v4 API.
  • pulsar-io-core now depends on the V5 API, so pulsar-client-api-v5, pulsar-client-v5,
    pulsar-tls-factory-api and pulsar-http-client-api are added to the NAR platform exclusions: the runtime
    provides them.
  • cleanupSubscription deletes a topic:// input's subscription through the scalable topics admin API.
  • Triggering a function with a topic:// topic is rejected with 400.

V5 client fix. The commit "[fix][client] Retry a V5 stream consumer's initial subscribe with the
current assignment" fixes a bug that this work exposed. When two stream consumers join a subscription at the
same moment, the first one's initial subscribe retried with a stale assignment and never completed. It can be
split into its own PR.

Not supported with the V5 client in this PR (each is rejected with an explicit error):

  • consumer and producer encryption;
  • message payload processors and consumer properties;
  • skipToLatest;
  • SinkContext seek, pause and resume (they throw UnsupportedOperationException at call time);
  • per-message schemas on one producer.

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • Unit tests:
    • ClientApiResolverTest: the selection matrix and every rejection.
    • FunctionConfigUtilsTest, SinkConfigUtilsTest and SourceConfigUtilsTest: round-trip, update and
      rejection of clientApi.
    • CmdFunctionsTest: --client-api.
    • V5InteropTest, V5ProducerAdapterTest, V5ProducerFactoryTest and LazyPulsarClientV5Test.
    • StreamAckTrackerTest, including out-of-order completion and cumulative acknowledgment.
    • ContextImplTest: getPulsarClientV5(), publishing to topic:// from v4 and V5 components, and the
      rejection of seek/pause/resume.
  • End-to-end, against a TLS and TLS-authentication broker with the function worker (thread runtime):
    • PulsarSourceV5E2ETest:
      • a source publishes to a topic:// topic;
      • with clientApi=V5 on a persistent:// output, the publisher carries the V5 marker that PIP-475
        migration checks;
      • V4 with a topic:// output is rejected.
    • PulsarFunctionV5E2ETest:
      • A Shared function on topic:// in and out; every input message is acknowledged, and deleting the
        function removes the subscription.
      • A Key_Shared function with parallelism 2; each key's messages arrive in order. Before the client fix,
        one instance never finished subscribing.
      • A clientApi=V5 Shared function on persistent:// topics; its consumer carries the V5 marker, and
        migrateToScalable(input, force=false) succeeds while the function is running.
  • Mutation checks:
    • forcing the runtime back to the v4 client fails both positive source tests;
    • dropping the scalable-topic branch of the subscription cleanup fails the Shared function test.
  • Regression:
    • the existing PulsarFunctionE2ETest, PulsarSinkE2ETest, PulsarSourceE2ETest,
      PulsarBatchSourceE2ETest and PulsarFunctionTlsTest pass unchanged;
    • all tests of the functions worker, runtime, runtime-all, instance and utils modules pass;
    • ./gradlew sanityCheck passes.
  • Not verified yet:
    • the process and Kubernetes runtimes, where user code loads java-instance.jar in its own class loader;
    • the batch source with V5 on a live broker;
    • the fatal-handler path of a failed stream record. It has no automated test, because building a real V5
      message in a unit test needs package-private client classes.

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)
    pulsar-functions-api now depends on pulsar-client-api-v5, and pulsar-functions-instance on
    pulsar-client-v5. Both are Pulsar modules, with no new third-party libraries.
  • The public API
    Added BaseContext.getPulsarClientV5(), a default method, and a clientApi field to FunctionConfig,
    SourceConfig and SinkConfig. Context.newOutputMessage now publishes to topic:// topics with the V5
    client.
  • The schema
  • The default values of configurations
  • The threading model
  • The binary protocol
  • The REST endpoints
    The function, source and sink configs accept clientApi. Trigger rejects topic:// topics.
  • The admin CLI options
    --client-api on pulsar-admin functions|sources|sinks create|update|localrun.
  • The metrics
  • Anything that affects deployment
    java-instance.jar now contains the V5 client API. NARs exclude the V5 client modules that the runtime
    provides.

…s and IO

Add a clientApi setting (V4/V5, unset = pick from the topic domains) to
FunctionConfig, SourceConfig and SinkConfig, FunctionDetails proto field 24
and --client-api on pulsar-admin functions/sources/sinks. ClientApiResolver
applies the pulsar-client/pulsar-perf selection rule plus the V5 runtime's
limits when a component is converted and again at instance startup, where
V5 is still rejected until the runtime is wired.

Add V5Interop, an internal accessor for the V5 client's schema and message
adapters.

Assisted-by: Claude Code (claude-opus-5-5)
A component whose topics resolve to the V5 client now publishes its output,
its context messages and its log topic with the V5 client. The runtime
factory owns one V5 client, created on first use and closed with the
factory, and passes it to the instance.

V5 producers are presented through the v4 Producer and TypedMessageBuilder
interfaces (V5ProducerAdapter, V5TypedMessageBuilder, V5MessageIdAdapter),
so PulsarSink, ProducerCache and ContextImpl keep a single code path;
V5ProducerFactory maps ProducerConfig onto the V5 producer builder.
InstanceUtils.createPulsarClientV5Builder applies the runtime's service
URL, authentication, TLS and memory settings to the V5 client.

A v4 component that publishes to a topic:// topic through its context now
fails with a hint to set clientApi to V5. Consuming with the V5 client is
still rejected at instance startup.

Assisted-by: Claude Code (claude-opus-5-5)
… current assignment

When two stream consumers join a subscription at the same time, the second
join rebalances the first consumer's initial assignment while it is still
subscribing. The first consumer's subscribe is then rejected with
ConsumerBusy for the segment that moved, and its retry used the assignment
captured at the first attempt, so it kept asking for that segment until the
operation timeout.

The assignment listener is only registered once the initial subscribe
succeeds, so updates that arrive in the meantime are only held by the
session. The retry now reads the session's current assignment.

Found by a Pulsar Function with parallelism 2 on a Key_Shared (stream)
subscription: both instances join at once, and before this change one of
them never finished subscribing.

Assisted-by: Claude Code (claude-opus-5-5)
A function or sink whose topics resolve to the V5 client now reads its
inputs with V5PulsarSource. Each input topic gets its own V5 consumer:

- a Shared subscription uses a QueueConsumer: individual acks, negative
  acks, the processing timeout, the negative-ack delay and the dead letter
  policy map from the component's settings;
- Failover, Key_Shared and EFFECTIVELY_ONCE use a StreamConsumer, which
  delivers in order and splits key ranges across the component's
  instances. A stream only acknowledges cumulatively, so StreamAckTracker
  acknowledges the completed prefix of the records in receive order; a
  failed record fails the instance, as EFFECTIVELY_ONCE already does.

Records are PulsarRecords over the v4 message that the V5 message wraps,
so connectors that read the schema, the message or the encryption
context keep working, and the record's topic is the topic:// input, not
the segment. Each instance uses a stable consumer name so that a restarted
instance rejoins its stream consumer group.

SinkContext seek, pause and resume throw UnsupportedOperationException in
V5 components. The resolver also counts inputs listed only in the
deprecated topicsToSerDeClassName map.

Assisted-by: Claude Code (claude-opus-5-5)
Add V5 accessors to the Functions and IO API, all default methods so that
other Context implementations keep compiling:

- BaseContext.getPulsarClientV5() and getPulsarClientBuilderV5();
- Context and SourceContext newOutputMessageV5(topic, schema), which
  publishes through the runtime's producer cache;
- Record.getMessageV5(), set for records read with the V5 client.

Every component can use them, whichever client its own topics use: the
runtime creates its shared V5 client on first use.

Packaging: pulsar-functions-api depends on pulsar-client-api-v5, so the V5
API is in java-instance.jar; the OpenTelemetry API that its TLS factory
SPI brings is kept out of java-instance.jar, as the v4 API does. NARs no
longer bundle the V5 client modules, which the runtime provides.

Worker: cleanupSubscription deletes a topic:// input's subscription with
the scalable topics admin API, and triggering a function with a topic://
input or output is rejected with 400, since triggering uses the worker's
v4 client.

Assisted-by: Claude Code (claude-opus-5-5)
- A failed record on a V5 stream subscription now fails the instance
  through its fatal handler. Record.fail() may be called from a producer
  callback or a connector's thread, where the thrown exception was lost
  and the acknowledgments stalled for good.
- A cumulative acknowledgment of a stream record completes every record
  received before it, as it does for v4 Failover subscriptions.
- The resolver rejects at submission what the V5 runtime cannot honor,
  instead of failing when the instance starts: consumer and producer
  encryption, message payload processors, consumer properties and
  skipToLatest, and maxMessageRetries, deadLetterTopic and timeoutMs with
  ordered subscriptions, which a stream would silently ignore.
- Without a configured memory limit the V5 client keeps its default one:
  the V5 producer has no pending-message limit, so it would otherwise hold
  an unbounded number of messages. Ignored producer settings are logged as
  warnings.
- The V5 client implementation is looked up with the runtime's
  classloader, so that the first V5 call from user code in the process and
  Kubernetes runtimes does not fail the one-time lookup.
- Smaller API: Context.newOutputMessage publishes to topic:// topics with
  the V5 client in every component, which replaces newOutputMessageV5;
  getPulsarClientBuilderV5 and Record.getMessageV5 are removed. The only
  API addition is BaseContext.getPulsarClientV5().

Assisted-by: Claude Code (claude-opus-5-5)
…r names

- When maxMessageRetries is set without a deadLetterTopic, the V5 queue
  consumer now uses the v4 client's default, <input>-<subscription>-DLQ in
  the input topic's domain, instead of the V5 client's default, a
  topic:// topic without the subscription name that has to exist
  beforehand. Switching a Shared component to the V5 client keeps its dead
  letter topic.
- Every start of an instance uses a new V5 consumer name. A stream consumer
  group treats a second consumer with the same name as a reconnect of the
  first, so a stable name could let two live instances share an
  assignment, and a clean stop already leaves the group explicitly.
- Producer settings that the V5 producer ignores are warned about only when
  they are set: the stored producer spec fills in 0 for unset limits.
- Correct the stated reason for keeping OpenTelemetry out of
  java-instance.jar.

Assisted-by: Claude Code (claude-opus-5-5)
…lient

Assisted-by: Claude Code (claude-opus-5-5)
…l the initial subscribe completes

Retrying the initial subscribe with the session's current assignment can
release a segment that the first attempt already attached. The release
drains the segment, waiting until every message handed out is acked, but
the segment's receive loop had already queued messages that the
application cannot ack: it gets the consumer only once the subscribe
completes. The drain then never finished and the subscribe hung.

The receive loops of the segments subscribed during the initial subscribe
now start once it completes, so nothing is handed out before the
application can ack it and a release during a retry drains at once.

Assisted-by: Claude Code (claude-opus-5-5)
- Delete the scalable topic's subscription when cleaning up a clientApi=V5
  component whose persistent:// input PIP-475 migrated: the V5 client
  follows the migrated topic under its persistent:// name, and deleting
  that topic's subscription left the new segments' behind.
- Make V5MessageIdAdapter Java-serializable through its V5 byte form; the
  v4 MessageId interface is Serializable, and the transient delegate
  deserialized as null.
- Log and count a publish error before failing the source record: failing
  a record of a V5 stream subscription throws, which skipped both.
- Close a V5 consumer that finishes subscribing after the source was
  closed, so that a stopped instance does not stay in the stream consumer
  group holding its key ranges.
- Take the ContextImpl lock only to create the producer factory for
  topic:// topics, not on every v4 publish.
- Warn about ignored producer settings once per producer factory.
- Import V5 types by simple name where no v4 type shares it.

Assisted-by: Claude Code (claude-opus-5-5)
@lhotari
lhotari force-pushed the lh-pip-496-fn-v5-client-impl branch from 7b37563 to f45f9ac Compare September 30, 2026 01:22
…ient-impl

# Conflicts:
#	pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/ComponentImpl.java
@lhotari
lhotari merged commit 99afb59 into apache:master Sep 30, 2026
43 checks passed
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.

2 participants