From dfd1ed51a251d45109ca323651ed70b6c4d6daa8 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 3 Jan 2024 00:17:01 +0800 Subject: [PATCH 1/2] [improve] [broker] PIP-299-part-1: Add dynamic config support: dispatcherPauseOnAckStatePersistentEnabled --- .../pulsar/broker/service/AbstractTopic.java | 8 ++- .../pulsar/broker/service/BrokerService.java | 26 ++++++++ ...PersistentDispatcherMultipleConsumers.java | 40 +++++++----- ...SubscriptionPauseOnAckStatPersistTest.java | 65 +++++++++++++++++++ 4 files changed, 122 insertions(+), 17 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java index ba35c8a280e9e..470f113e369bc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java @@ -397,7 +397,8 @@ private void updateTopicPolicyByBrokerConfig() { topicPolicies.getSchemaValidationEnforced().updateBrokerValue(config.isSchemaValidationEnforced()); topicPolicies.getEntryFilters().updateBrokerValue(new EntryFilters(String.join(",", config.getEntryFilterNames()))); - + topicPolicies.getDispatcherPauseOnAckStatePersistentEnabled() + .updateBrokerValue(config.isDispatcherPauseOnAckStatePersistentEnabled()); updateEntryFilters(); } @@ -1267,6 +1268,11 @@ public void updateBrokerDispatchRate() { dispatchRateInBroker(brokerService.pulsar().getConfiguration())); } + public void updateDispatchPauseOnAckStatePersistentEnabled() { + topicPolicies.getDispatcherPauseOnAckStatePersistentEnabled().updateBrokerValue( + brokerService.pulsar().getConfiguration().isDispatcherPauseOnAckStatePersistentEnabled()); + } + public void addFilteredEntriesCount(int filtered) { this.filteredEntriesCounter.add(filtered); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 4077762bb0640..d24058bc85b77 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -2605,6 +2605,10 @@ private void updateConfigurationAndRegisterListeners() { registerConfigurationListener("dispatchThrottlingRatePerSubscriptionInByte", (dispatchRatePerTopicInByte) -> { updateSubscriptionMessageDispatchRate(); }); + // add listener to update message-dispatch-rate in byte for subscription + registerConfigurationListener("dispatcherPauseOnAckStatePersistentEnabled", (dispatchRatePerTopicInByte) -> { + updateDispatchPauseOnAckStatePersistentEnabled(); + }); // add listener to update message-dispatch-rate in msg for replicator registerConfigurationListener("dispatchThrottlingRatePerReplicatorInMsg", @@ -2743,6 +2747,28 @@ private void updateTopicMessageDispatchRate() { }); } + private void updateDispatchPauseOnAckStatePersistentEnabled() { + this.pulsar().getExecutor().execute(() -> { + forEachTopic(topic -> { + if (topic instanceof PersistentTopic) { + // Update policies. + PersistentTopic persistentTopic = (PersistentTopic) topic; + persistentTopic.updateDispatchPauseOnAckStatePersistentEnabled(); + // Trigger new read if subscriptions has been paused before. + if (!pulsar().getConfiguration().isDispatcherPauseOnAckStatePersistentEnabled()) { + persistentTopic.updateDispatchPauseOnAckStatePersistentEnabled(); + persistentTopic.getSubscriptions().forEach((sName, subscription) -> { + if (subscription.getDispatcher() == null) { + return; + } + subscription.getDispatcher().afterAckMessages(null, 0); + }); + } + } + }); + }); + } + private void updateBrokerSubscriptionTypesEnabled(Object subscriptionTypesEnabled) { this.pulsar().getExecutor().execute(() -> { // update subscriptionTypesEnabled diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 0d1f198a7ca7e..56e9fe52df678 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -1039,23 +1039,31 @@ public void addUnAckedMessages(int numberOfMessages) { @Override public void afterAckMessages(Throwable exOfDeletion, Object ctxOfDeletion) { - if (blockedDispatcherOnCursorDataCanNotFullyPersist == TRUE) { - if (cursor.isCursorDataFullyPersistable()) { - // If there was no previous pause due to cursor data is too large to persist, we don't need to manually - // trigger a new read. This can avoid too many CPU circles. - if (BLOCKED_DISPATCHER_ON_CURSOR_DATA_CAN_NOT_FULLY_PERSIST_UPDATER.compareAndSet(this, TRUE, FALSE)) { - readMoreEntriesAsync(); - } else { - // Retry due to conflict update. - afterAckMessages(exOfDeletion, ctxOfDeletion); - } + boolean paused = blockedDispatcherOnCursorDataCanNotFullyPersist == TRUE; + boolean shouldPauseNow = !cursor.isCursorDataFullyPersistable() + && topic.isDispatcherPauseOnAckStatePersistentEnabled(); + // No need to change. + if (paused == shouldPauseNow) { + return; + } + // Should change to "un-pause". + if (paused && !shouldPauseNow) { + // If there was no previous pause due to cursor data is too large to persist, we don't need to manually + // trigger a new read. This can avoid too many CPU circles. + if (BLOCKED_DISPATCHER_ON_CURSOR_DATA_CAN_NOT_FULLY_PERSIST_UPDATER.compareAndSet(this, TRUE, FALSE)) { + readMoreEntriesAsync(); + } else { + // Retry due to conflict update. + afterAckMessages(exOfDeletion, ctxOfDeletion); } - } else { - if (!cursor.isCursorDataFullyPersistable()) { - if (BLOCKED_DISPATCHER_ON_CURSOR_DATA_CAN_NOT_FULLY_PERSIST_UPDATER.compareAndSet(this, FALSE, TRUE)) { - // Retry due to conflict update. - afterAckMessages(exOfDeletion, ctxOfDeletion); - } + return; + } + // Should change to "paused". + if (!paused && shouldPauseNow) { + if (!BLOCKED_DISPATCHER_ON_CURSOR_DATA_CAN_NOT_FULLY_PERSIST_UPDATER + .compareAndSet(this, FALSE, TRUE)) { + // Retry due to conflict update. + afterAckMessages(exOfDeletion, ctxOfDeletion); } } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionPauseOnAckStatPersistTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionPauseOnAckStatPersistTest.java index 0029e61df4c49..02a906e91a38b 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionPauseOnAckStatPersistTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionPauseOnAckStatPersistTest.java @@ -37,7 +37,9 @@ import org.apache.pulsar.client.admin.GetStatsOptions; import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.HierarchyTopicPolicies; import org.apache.pulsar.common.policies.data.TopicPolicies; +import org.awaitility.Awaitility; import org.awaitility.reflect.WhiteboxImpl; import org.testng.Assert; import org.testng.annotations.AfterClass; @@ -211,6 +213,69 @@ public boolean hasAckedMessage(String v) { } } + @Test + public void testBrokerDynamicConfig() throws Exception { + final String tpName = BrokerTestUtil.newUniqueName("persistent://public/default/tp"); + final String subscription = "s1"; + final int msgSendCount = MAX_UNACKED_RANGES_TO_PERSIST * 4; + final int incomingQueueSize = MAX_UNACKED_RANGES_TO_PERSIST * 10; + + // Enable "dispatcherPauseOnAckStatePersistentEnabled". + admin.brokers().updateDynamicConfiguration("dispatcherPauseOnAckStatePersistentEnabled", "true"); + admin.topics().createNonPartitionedTopic(tpName); + admin.topics().createSubscription(tpName, subscription, MessageId.earliest); + + PersistentTopic persistentTopic = + (PersistentTopic) pulsar.getBrokerService().getTopic(tpName, false).join().get(); + Awaitility.await().untilAsserted(() -> { + Assert.assertTrue(pulsar.getConfig().isDispatcherPauseOnAckStatePersistentEnabled()); + HierarchyTopicPolicies policies = WhiteboxImpl.getInternalState(persistentTopic, "topicPolicies"); + Boolean v = policies.getDispatcherPauseOnAckStatePersistentEnabled().get(); + Assert.assertNotNull(v); + Assert.assertTrue(v.booleanValue()); + }); + + // Send double MAX_UNACKED_RANGES_TO_PERSIST messages. + Producer p1 = pulsarClient.newProducer(Schema.STRING).topic(tpName).enableBatching(false).create(); + ArrayList messageIdsSent = new ArrayList<>(); + for (int i = 0; i < msgSendCount; i++) { + MessageIdImpl messageId = (MessageIdImpl) p1.send(Integer.valueOf(i).toString()); + messageIdsSent.add(messageId); + } + // Make ack holes. + Consumer c1 = pulsarClient.newConsumer(Schema.STRING).topic(tpName).subscriptionName(subscription) + .receiverQueueSize(incomingQueueSize).isAckReceiptEnabled(true) + .subscriptionType(SubscriptionType.Shared).subscribe(); + ackOddMessagesOnly(c1); + + cancelPendingRead(tpName, subscription); + triggerNewReadMoreEntries(tpName, subscription); + + // Verify: the dispatcher has been paused. + final String specifiedMessage = "9876543210"; + p1.send(specifiedMessage); + Message msg1 = c1.receive(2, TimeUnit.SECONDS); + Assert.assertNull(msg1, msg1 == null ? "null" : msg1.getValue()); + + // Disable "dispatcherPauseOnAckStatePersistentEnabled". + admin.brokers().updateDynamicConfiguration("dispatcherPauseOnAckStatePersistentEnabled", "false"); + Awaitility.await().untilAsserted(() -> { + Assert.assertFalse(pulsar.getConfig().isDispatcherPauseOnAckStatePersistentEnabled()); + HierarchyTopicPolicies policies = WhiteboxImpl.getInternalState(persistentTopic, "topicPolicies"); + Boolean v = policies.getDispatcherPauseOnAckStatePersistentEnabled().get(); + Assert.assertTrue(v == null || !v.booleanValue()); + }); + + // Verify the new message can be received. + Message msg2 = c1.receive(2, TimeUnit.SECONDS); + Assert.assertNotNull(msg2); + Assert.assertEquals(msg2.getValue(), specifiedMessage); + // cleanup. + p1.close(); + c1.close(); + admin.topics().delete(tpName, false); + } + @Test(dataProvider = "multiConsumerSubscriptionTypes") public void testPauseOnAckStatPersist(SubscriptionType subscriptionType) throws Exception { final String tpName = BrokerTestUtil.newUniqueName("persistent://public/default/tp"); From f455b8d3db451759ae6a404a45786c453cf95a48 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 4 Jan 2024 19:44:23 +0800 Subject: [PATCH 2/2] [fix] [client] Messages lost due to TopicListWatcher reconnect --- .../auth/MockedPulsarServiceBaseTest.java | 8 ++ .../impl/PatternTopicsConsumerImplTest.java | 63 ++++++++++++--- .../impl/PatternMultiTopicsConsumerImpl.java | 79 ++++++++++++++++--- .../pulsar/client/impl/TopicListWatcher.java | 7 +- .../client/impl/TopicListWatcherTest.java | 2 +- 5 files changed, 133 insertions(+), 26 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java index b8d75bd0fbcac..cc5ea3bbb7bd4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java @@ -708,5 +708,13 @@ public static class ServiceProducer { private PersistentTopic persistentTopic; } + protected void sleepSeconds(int seconds){ + try { + Thread.currentThread().sleep(1000 * seconds); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } + } + private static final Logger log = LoggerFactory.getLogger(MockedPulsarServiceBaseTest.class); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/PatternTopicsConsumerImplTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/PatternTopicsConsumerImplTest.java index 451f93067b2ca..9115cefa6e1e6 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/PatternTopicsConsumerImplTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/PatternTopicsConsumerImplTest.java @@ -37,14 +37,18 @@ import io.netty.util.Timeout; import org.apache.pulsar.broker.namespace.NamespaceService; import org.apache.pulsar.client.api.Consumer; +import org.apache.pulsar.client.api.InjectedClientCnxClientBuilder; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageRoutingMode; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.ProducerConsumerBase; +import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.RegexSubscriptionMode; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; +import org.apache.pulsar.common.api.proto.BaseCommand; +import org.apache.pulsar.common.api.proto.CommandWatchTopicListSuccess; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.policies.data.TenantInfoImpl; import org.awaitility.Awaitility; @@ -53,6 +57,7 @@ import org.slf4j.LoggerFactory; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; +import org.testng.annotations.DataProvider; import org.testng.annotations.Test; @Test(groups = "broker-impl") @@ -620,13 +625,28 @@ public void testStartEmptyPatternConsumer() throws Exception { producer3.close(); } - @Test(timeOut = testTimeout) - public void testAutoSubscribePatterConsumerFromBrokerWatcher() throws Exception { - String key = "AutoSubscribePatternConsumer"; - String subscriptionName = "my-ex-subscription-" + key; + @DataProvider(name= "delayTypesOfWatchingTopics") + public Object[][] delayTypesOfWatchingTopics(){ + return new Object[][]{ + {true}, + {false} + }; + } - Pattern pattern = Pattern.compile("persistent://my-property/my-ns/pattern-topic.*"); - Consumer consumer = pulsarClient.newConsumer() + @Test(timeOut = testTimeout, dataProvider = "delayTypesOfWatchingTopics") + public void testAutoSubscribePatterConsumerFromBrokerWatcher(boolean delayWatchingTopics) throws Exception { + final String key = "AutoSubscribePatternConsumer"; + final String subscriptionName = "my-ex-subscription-" + key; + final Pattern pattern = Pattern.compile("persistent://my-property/my-ns/pattern-topic.*"); + + PulsarClient client = null; + if (delayWatchingTopics) { + client = createDelayWatchTopicsClient(); + } else { + client = pulsarClient; + } + + Consumer consumer = client.newConsumer() .topicsPattern(pattern) // Disable automatic discovery. .patternAutoDiscoveryPeriod(1000) @@ -636,12 +656,6 @@ public void testAutoSubscribePatterConsumerFromBrokerWatcher() throws Exception .receiverQueueSize(4) .subscribe(); - // Wait topic list watcher creation. - Awaitility.await().untilAsserted(() -> { - CompletableFuture completableFuture = WhiteboxImpl.getInternalState(consumer, "watcherFuture"); - assertTrue(completableFuture.isDone() && !completableFuture.isCompletedExceptionally()); - }); - // 1. create partition String topicName = "persistent://my-property/my-ns/pattern-topic-1-" + key; TenantInfoImpl tenantInfo = createDefaultTenantInfo(); @@ -657,7 +671,32 @@ public void testAutoSubscribePatterConsumerFromBrokerWatcher() throws Exception assertEquals(((PatternMultiTopicsConsumerImpl) consumer).getPartitionedTopics().size(), 1); }); + // cleanup. consumer.close(); + admin.topics().deletePartitionedTopic(topicName); + } + + private PulsarClient createDelayWatchTopicsClient() throws Exception { + ClientBuilderImpl clientBuilder = (ClientBuilderImpl) PulsarClient.builder().serviceUrl(lookupUrl.toString()); + return InjectedClientCnxClientBuilder.create(clientBuilder, + (conf, eventLoopGroup) -> new ClientCnx(conf, eventLoopGroup) { + public CompletableFuture newWatchTopicList( + BaseCommand command, long requestId) { + // Inject 2 seconds delay when sending command New Watch Topics. + CompletableFuture res = new CompletableFuture<>(); + new Thread(() -> { + sleepSeconds(2); + super.newWatchTopicList(command, requestId).whenComplete((v, ex) -> { + if (ex != null) { + res.completeExceptionally(ex); + } else { + res.complete(v); + } + }); + }).start(); + return res; + } + }); } // simulate subscribe a pattern which has 3 topics, but then matched topic added in. diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PatternMultiTopicsConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PatternMultiTopicsConsumerImpl.java index c6ea6216cc1f4..ca6895c333b8e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PatternMultiTopicsConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PatternMultiTopicsConsumerImpl.java @@ -31,6 +31,7 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; import java.util.regex.Pattern; import java.util.stream.Collectors; import org.apache.pulsar.client.api.Consumer; @@ -50,9 +51,12 @@ public class PatternMultiTopicsConsumerImpl extends MultiTopicsConsumerImpl watcherFuture; + private final CompletableFuture watcherFuture = new CompletableFuture<>(); protected NamespaceName namespaceName; private volatile Timeout recheckPatternTimeout = null; + private volatile AtomicReference retryRecheckPatternTask = new AtomicReference<>(); + private final Backoff retryRecheckPatternTaskBackoff = new BackoffBuilder().setInitialTime(5, TimeUnit.SECONDS) + .setMax(1, TimeUnit.MINUTES).setMandatoryStop(0, TimeUnit.SECONDS).create(); private volatile String topicsHash; public PatternMultiTopicsConsumerImpl(Pattern topicsPattern, @@ -78,11 +82,10 @@ public PatternMultiTopicsConsumerImpl(Pattern topicsPattern, this.topicsChangeListener = new PatternTopicsChangedListener(); this.recheckPatternTimeout = client.timer() .newTimeout(this, Math.max(1, conf.getPatternAutoDiscoveryPeriod()), TimeUnit.SECONDS); - this.watcherFuture = new CompletableFuture<>(); if (subscriptionMode == Mode.PERSISTENT) { long watcherId = client.newTopicListWatcherId(); new TopicListWatcher(topicsChangeListener, client, topicsPattern, watcherId, - namespaceName, topicsHash, watcherFuture); + namespaceName, topicsHash, watcherFuture, () -> recheckTopicsChangeRetryIfFailed()); watcherFuture .thenAccept(__ -> recheckPatternTimeout.cancel()) .exceptionally(ex -> { @@ -105,7 +108,62 @@ public void run(Timeout timeout) throws Exception { if (timeout.isCancelled()) { return; } - client.getLookup().getTopicsUnderNamespace(namespaceName, subscriptionMode, topicsPattern.pattern(), topicsHash) + asyncRecheckTopicsChange().exceptionally(ex -> { + log.warn("[{}] Failed to recheck topics change: {}", topic, ex.getMessage()); + return null; + }).thenAccept(__ -> { + // schedule the next re-check task + this.recheckPatternTimeout = client.timer() + .newTimeout(PatternMultiTopicsConsumerImpl.this, + Math.max(1, conf.getPatternAutoDiscoveryPeriod()), TimeUnit.SECONDS); + }); + } + + private void recheckTopicsChangeRetryIfFailed() { + recheckTopicsChangeRetryIfFailed(null); + } + + private void recheckTopicsChangeRetryIfFailed(Timeout retryTask) { + // This method will be called by A New Call or Timeout scheduled call. + final boolean isNew = (retryTask == null); + // Skip if closed or the task has been cancelled. + if (getState() == State.Closing || getState() == State.Closed + || (retryTask != null && retryTask.isCancelled())) { + retryRecheckPatternTask.compareAndSet(retryTask, null); + return; + } + // Skip the new check if contains a retry task. + Timeout pendingRetryTask = retryRecheckPatternTask.get(); + if (isNew && pendingRetryTask != null) { + return; + } + // Do check. + asyncRecheckTopicsChange().whenComplete((ignore, ex) -> { + if (ex != null) { + log.warn("[{}] Failed to recheck topics change: {}", topic, ex.getMessage()); + long delayMs = retryRecheckPatternTaskBackoff.next(); + Timeout newTask = client.timer().newTimeout(timeout -> { + if (timeout.cancel()) { + return; + } + recheckTopicsChangeRetryIfFailed(); + }, delayMs, TimeUnit.MILLISECONDS); + if (!retryRecheckPatternTask.compareAndSet(retryTask, newTask)) { + // Another thread added a new task, so cancel current one. + newTask.cancel(); + } + } else { + retryRecheckPatternTaskBackoff.reset(); + if (!isNew) { + retryRecheckPatternTask.compareAndSet(retryTask, null); + } + } + }); + } + + private CompletableFuture asyncRecheckTopicsChange() { + String pattern = topicsPattern.pattern(); + return client.getLookup().getTopicsUnderNamespace(namespaceName, subscriptionMode, pattern, topicsHash) .thenCompose(getTopicsResult -> { if (log.isDebugEnabled()) { @@ -125,14 +183,6 @@ public void run(Timeout timeout) throws Exception { } return updateSubscriptions(topicsPattern, this::setTopicsHash, getTopicsResult, topicsChangeListener, oldTopics); - }).exceptionally(ex -> { - log.warn("[{}] Failed to recheck topics change: {}", topic, ex.getMessage()); - return null; - }).thenAccept(__ -> { - // schedule the next re-check task - this.recheckPatternTimeout = client.timer() - .newTimeout(PatternMultiTopicsConsumerImpl.this, - Math.max(1, conf.getPatternAutoDiscoveryPeriod()), TimeUnit.SECONDS); }); } @@ -234,6 +284,11 @@ public CompletableFuture closeAsync() { timeout.cancel(); recheckPatternTimeout = null; } + Timeout retryTaskToRecheckTopics = retryRecheckPatternTask.get(); + if (retryTaskToRecheckTopics != null) { + retryTaskToRecheckTopics.cancel(); + retryRecheckPatternTask.compareAndSet(retryTaskToRecheckTopics, null); + } List> closeFutures = new ArrayList<>(2); if (watcherFuture.isDone() && !watcherFuture.isCompletedExceptionally()) { TopicListWatcher watcher = watcherFuture.getNow(null); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/TopicListWatcher.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/TopicListWatcher.java index 2ce784dbaac04..489a07a606eb2 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/TopicListWatcher.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/TopicListWatcher.java @@ -56,11 +56,14 @@ public class TopicListWatcher extends HandlerState implements ConnectionHandler. private final List previousExceptions = new CopyOnWriteArrayList<>(); private final AtomicReference clientCnxUsedForWatcherRegistration = new AtomicReference<>(); + private final Runnable recheckTopicsChangeAfterReconnect; + public TopicListWatcher(PatternMultiTopicsConsumerImpl.TopicsChangedListener topicsChangeListener, PulsarClientImpl client, Pattern topicsPattern, long watcherId, NamespaceName namespace, String topicsHash, - CompletableFuture watcherFuture) { + CompletableFuture watcherFuture, + Runnable recheckTopicsChangeAfterReconnect) { super(client, topicsPattern.pattern()); this.topicsChangeListener = topicsChangeListener; this.name = "Watcher(" + topicsPattern + ")"; @@ -77,6 +80,7 @@ public TopicListWatcher(PatternMultiTopicsConsumerImpl.TopicsChangedListener top this.namespace = namespace; this.topicsHash = topicsHash; this.watcherFuture = watcherFuture; + this.recheckTopicsChangeAfterReconnect = recheckTopicsChangeAfterReconnect; connectionHandler.grabCnx(); } @@ -141,6 +145,7 @@ public CompletableFuture connectionOpened(ClientCnx cnx) { this.connectionHandler.resetBackoff(); + recheckTopicsChangeAfterReconnect.run(); watcherFuture.complete(this); future.complete(null); }).exceptionally((e) -> { diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/TopicListWatcherTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/TopicListWatcherTest.java index dd75770b5688d..7e9fd601d4f67 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/TopicListWatcherTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/TopicListWatcherTest.java @@ -71,7 +71,7 @@ public void setup() { watcherFuture = new CompletableFuture<>(); watcher = new TopicListWatcher(listener, client, Pattern.compile(topic), 7, - NamespaceName.get("tenant/ns"), null, watcherFuture); + NamespaceName.get("tenant/ns"), null, watcherFuture, () -> {}); } @Test