[fix][client] The input parameters for PulsarPullConsumerImpl have been changed from PulsarAdmin parameters to Supplier format, requiring dynamic updates. (#27)
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 @@
private final Map<String, Consumer<T>> consumerMap;
private final OffsetToMessageIdCache offsetToMessageIdCache;
private final ReaderCache<T> readerCache;
- private final PulsarAdmin pulsarAdmin;
+ private final Supplier<PulsarAdmin> pulsarAdminSupplier;
private final Supplier<PulsarClient> pulsarClientSupplier;
private final ConsumerBuilder<T> consumerBuilder;
@@ -74,7 +74,7 @@
String brokerCluster,
Schema<T> schema,
Supplier<PulsarClient> clientSupplier,
- PulsarAdmin admin,
+ Supplier<PulsarAdmin> adminSupplier,
ConsumerBuilder<T> 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 @@
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 @@
}
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 @@
"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<T> pull(PullRequest request) {
validatePullParameters(request.getMaxMessages(), request.getMaxBytes());
@@ -231,14 +239,15 @@
@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 @@
brokerCluster,
Schema.BYTES,
() -> pulsarClient,
- pulsarAdmin,
+ () -> pulsarAdmin,
null);
pullConsumer.start();
@@ -177,7 +177,7 @@
brokerCluster,
Schema.BYTES,
() -> pulsarClient,
- pulsarAdmin,
+ () -> pulsarAdmin,
null);
pullConsumer.start();
@@ -218,7 +218,7 @@
brokerCluster,
Schema.BYTES,
() -> pulsarClient,
- pulsarAdmin,
+ () -> pulsarAdmin,
null);
pullConsumer.start();