[feat][fn] PIP-496: Support the V5 client and topic:// topics in Pulsar Functions and IO - #26697
Merged
Merged
Conversation
…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
force-pushed
the
lh-pip-496-fn-v5-client-impl
branch
from
September 30, 2026 01:22
7b37563 to
f45f9ac
Compare
…ient-impl # Conflicts: # pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/ComponentImpl.java
merlimat
approved these changes
Sep 30, 2026
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.
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 scalabletopic 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 bemigrated.
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-clientandpulsar-perfCLIs (#26693):topic://topic selects the V5 client, and any other domain selects v4.clientApi=V5makes the V5 client drivepersistent://topics.The new
clientApisetting (V4/V5; unset means "pick from the topics") is added toFunctionConfig,SourceConfigandSinkConfig, as proto fieldFunctionDetails.clientApi = 24(AUTOby default) and as--client-apionpulsar-admin functions|sources|sinks.ClientApiResolvervalidates it when a componentis converted, which is when the worker registers or updates it, and again when an instance starts. It rejects:
topic://with other domains, or naming asegment://topic directly;topic://topic, and V5 with anon-persistent://topic;sequence IDs are not monotonic across splits and merges);
persistent://inputs: any subscription type other than Shared, because a V5 stream consumercannot yet subscribe to a topic that has not been migrated;
skipToLatest;maxMessageRetries,deadLetterTopicandtimeoutMs, which a streamcannot 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.QueueConsumer. Individual acks, negative acks, the processing timeout, thenegative-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 V5keeps its dead letter topic.
StreamConsumer, which delivers in order and splits keyranges across the component's instances.
StreamAckTrackeracknowledges the completed prefix of thereceived records. Records completed out of order are never skipped.
from whichever thread calls
fail().<tenant>-<namespace>-<name>-<instanceId>-<random>: astream consumer group treats a reused name as a reconnect of the previous member.
PulsarRecords over the v4 message that the V5 message wraps. Connectors that read theschema, 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/MessageIdinterfaces, so
PulsarSink,ProducerCache,ContextImpl.newOutputMessageand the log appender keep asingle code path; only producer creation branches (
V5ProducerFactory).Context.newOutputMessagepublishesto a
topic://topic with the V5 client in every component.Runtime.
closed with the factory.
InstanceUtils.createPulsarClientV5Builderapplies the service URL, authentication, TLS and memory limitof 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.
pending-message limit, so the memory limit is what bounds it.
from user code in the process and Kubernetes runtimes cannot fail the JVM-wide lookup.
User API (PIP-496): one
defaultmethod,BaseContext.getPulsarClientV5(), which returns the runtime'sshared V5 client.
Packaging and worker.
pulsar-functions-apidepends onpulsar-client-api-v5, so the V5 API is injava-instance.jar. TheOpenTelemetry API is kept out of that jar, as with the v4 API.
pulsar-io-corenow depends on the V5 API, sopulsar-client-api-v5,pulsar-client-v5,pulsar-tls-factory-apiandpulsar-http-client-apiare added to the NAR platform exclusions: the runtimeprovides them.
cleanupSubscriptiondeletes atopic://input's subscription through the scalable topics admin API.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):
skipToLatest;seek,pauseandresume(they throwUnsupportedOperationExceptionat call time);Verifying this change
This change added tests and can be verified as follows:
ClientApiResolverTest: the selection matrix and every rejection.FunctionConfigUtilsTest,SinkConfigUtilsTestandSourceConfigUtilsTest: round-trip, update andrejection of
clientApi.CmdFunctionsTest:--client-api.V5InteropTest,V5ProducerAdapterTest,V5ProducerFactoryTestandLazyPulsarClientV5Test.StreamAckTrackerTest, including out-of-order completion and cumulative acknowledgment.ContextImplTest:getPulsarClientV5(), publishing totopic://from v4 and V5 components, and therejection of seek/pause/resume.
PulsarSourceV5E2ETest:topic://topic;clientApi=V5on apersistent://output, the publisher carries the V5 marker that PIP-475migration checks;
topic://output is rejected.PulsarFunctionV5E2ETest:topic://in and out; every input message is acknowledged, and deleting thefunction removes the subscription.
one instance never finished subscribing.
clientApi=V5Shared function onpersistent://topics; its consumer carries the V5 marker, andmigrateToScalable(input, force=false)succeeds while the function is running.PulsarFunctionE2ETest,PulsarSinkE2ETest,PulsarSourceE2ETest,PulsarBatchSourceE2ETestandPulsarFunctionTlsTestpass unchanged;worker,runtime,runtime-all,instanceandutilsmodules pass;./gradlew sanityCheckpasses.java-instance.jarin its own class loader;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
pulsar-functions-apinow depends onpulsar-client-api-v5, andpulsar-functions-instanceonpulsar-client-v5. Both are Pulsar modules, with no new third-party libraries.Added
BaseContext.getPulsarClientV5(), adefaultmethod, and aclientApifield toFunctionConfig,SourceConfigandSinkConfig.Context.newOutputMessagenow publishes totopic://topics with the V5client.
The function, source and sink configs accept
clientApi. Trigger rejectstopic://topics.--client-apionpulsar-admin functions|sources|sinks create|update|localrun.java-instance.jarnow contains the V5 client API. NARs exclude the V5 client modules that the runtimeprovides.