From eace249415ef10f2735a3c56eaa02981123e8ad6 Mon Sep 17 00:00:00 2001 From: nuolin Date: Wed, 5 Aug 2026 17:51:22 +0800 Subject: [PATCH] [fix][client] The input parameters for PulsarPullConsumerImpl have been changed from PulsarAdmin parameters to Supplier format, requiring dynamic updates. --- .../api/impl/PulsarPullConsumerImpl.java | 23 +++++++++++++------ .../src/test/java/PulsarPullConsumerTest.java | 6 ++--- 2 files changed, 19 insertions(+), 10 deletions(-) diff --git a/pulsar-client-common-contrib/src/main/java/org/apache/pulsar/client/api/impl/PulsarPullConsumerImpl.java b/pulsar-client-common-contrib/src/main/java/org/apache/pulsar/client/api/impl/PulsarPullConsumerImpl.java index 74bba7b..39a2175 100644 --- a/pulsar-client-common-contrib/src/main/java/org/apache/pulsar/client/api/impl/PulsarPullConsumerImpl.java +++ b/pulsar-client-common-contrib/src/main/java/org/apache/pulsar/client/api/impl/PulsarPullConsumerImpl.java @@ -62,7 +62,7 @@ public class PulsarPullConsumerImpl implements PulsarPullConsumer { private final Map> consumerMap; private final OffsetToMessageIdCache offsetToMessageIdCache; private final ReaderCache readerCache; - private final PulsarAdmin pulsarAdmin; + private final Supplier pulsarAdminSupplier; private final Supplier pulsarClientSupplier; private final ConsumerBuilder consumerBuilder; @@ -74,7 +74,7 @@ public PulsarPullConsumerImpl( String brokerCluster, Schema schema, Supplier clientSupplier, - PulsarAdmin admin, + Supplier adminSupplier, ConsumerBuilder consumerBuilder) { this.topic = Objects.requireNonNull(topic, "Topic must not be null"); this.subscription = Objects.requireNonNull(subscription, "Subscription must not be null"); @@ -82,10 +82,11 @@ public PulsarPullConsumerImpl( this.schema = Objects.requireNonNull(schema, "Schema must not be null"); this.pulsarClientSupplier = Objects.requireNonNull(clientSupplier, "PulsarClient must not be null"); - this.pulsarAdmin = Objects.requireNonNull(admin, "PulsarAdmin must not be null"); + this.pulsarAdminSupplier = + Objects.requireNonNull(adminSupplier, "PulsarAdmin must not be null"); this.consumerMap = new ConcurrentHashMap<>(); this.offsetToMessageIdCache = - OffsetToMessageIdCacheProvider.getOrCreateCache(admin, brokerCluster); + OffsetToMessageIdCacheProvider.getOrCreateCache(getPulsarAdmin(), brokerCluster); this.readerCache = ReaderCacheProvider.getOrCreateReaderCache( this.subscription, brokerCluster, schema, clientSupplier.get(), offsetToMessageIdCache); @@ -105,7 +106,8 @@ public void start() throws PulsarClientException { } private void initializePartitions() throws PulsarAdminException, PulsarClientException { - PartitionedTopicMetadata metadata = pulsarAdmin.topics().getPartitionedTopicMetadata(topic); + PartitionedTopicMetadata metadata = + getPulsarAdmin().topics().getPartitionedTopicMetadata(topic); this.partitionCount = metadata.partitions; if (partitionCount == 0) { @@ -132,6 +134,12 @@ private PulsarClient getPulsarClient() { "PulsarClient supplier returned null. Ensure PulsarClient is properly initialized."); } + private PulsarAdmin getPulsarAdmin() { + return Objects.requireNonNull( + pulsarAdminSupplier.get(), + "PulsarAdmin supplier returned null. Ensure PulsarAdmin is properly initialized."); + } + @Override public PullResponse pull(PullRequest request) { validatePullParameters(request.getMaxMessages(), request.getMaxBytes()); @@ -231,14 +239,15 @@ public void ack(long offset, int partition, CommandAck.AckType ackType) @Override public long searchOffset(int partition, long timestamp) throws PulsarAdminException { String partitionTopic = buildPartitionTopic(topic, partition); - return PulsarAdminUtils.searchOffset(partitionTopic, timestamp, brokerCluster, pulsarAdmin); + return PulsarAdminUtils.searchOffset( + partitionTopic, timestamp, brokerCluster, getPulsarAdmin()); } @Override public ConsumeStats getConsumeStats(int partition) throws PulsarAdminException { String partitionTopic = buildPartitionTopic(topic, partition); return PulsarAdminUtils.getConsumeStats( - partitionTopic, partition, subscription, brokerCluster, pulsarAdmin); + partitionTopic, partition, subscription, brokerCluster, getPulsarAdmin()); } @Override diff --git a/pulsar-client-common-contrib/src/test/java/PulsarPullConsumerTest.java b/pulsar-client-common-contrib/src/test/java/PulsarPullConsumerTest.java index dc49026..a2f5652 100644 --- a/pulsar-client-common-contrib/src/test/java/PulsarPullConsumerTest.java +++ b/pulsar-client-common-contrib/src/test/java/PulsarPullConsumerTest.java @@ -106,7 +106,7 @@ public void testPullConsumer(String topic, int partitionIndex) throws Exception brokerCluster, Schema.BYTES, () -> pulsarClient, - pulsarAdmin, + () -> pulsarAdmin, null); pullConsumer.start(); @@ -177,7 +177,7 @@ public void testSearchOffset(String topic, int partitionIndex) throws Exception brokerCluster, Schema.BYTES, () -> pulsarClient, - pulsarAdmin, + () -> pulsarAdmin, null); pullConsumer.start(); @@ -218,7 +218,7 @@ public void testGetConsumeStats() { brokerCluster, Schema.BYTES, () -> pulsarClient, - pulsarAdmin, + () -> pulsarAdmin, null); pullConsumer.start();