From 79c22199c96c02dab60403a9ecd86f172e43e143 Mon Sep 17 00:00:00 2001 From: technoboy Date: Sun, 6 Mar 2022 21:18:58 +0800 Subject: [PATCH 01/18] Fix update replication cluster but not update replicator --- .../service/persistent/PersistentTopic.java | 50 ++++++++++--------- .../pulsar/broker/service/ReplicatorTest.java | 29 +++++++++++ 2 files changed, 56 insertions(+), 23 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index a7c1714186974..7c292dae12bea 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -27,7 +27,6 @@ import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.Lists; import com.google.common.collect.Maps; -import com.google.common.collect.Sets; import io.netty.buffer.ByteBuf; import io.netty.util.concurrent.FastThreadLocal; import java.time.Clock; @@ -623,28 +622,13 @@ private boolean hasRemoteProducers() { public CompletableFuture startReplProducers() { // read repl-cluster from policies to avoid restart of replicator which are in process of disconnect and close - return brokerService.pulsar().getPulsarResources().getNamespaceResources() - .getPoliciesAsync(TopicName.get(topic).getNamespaceObject()) - .thenAccept(optPolicies -> { - if (optPolicies.isPresent()) { - if (optPolicies.get().replication_clusters != null) { - Set configuredClusters = Sets.newTreeSet(optPolicies.get().replication_clusters); - replicators.forEach((region, replicator) -> { - if (configuredClusters.contains(region)) { - replicator.startProducer(); - } - }); - } - } else { - replicators.forEach((region, replicator) -> replicator.startProducer()); - } - }).exceptionally(ex -> { - if (log.isDebugEnabled()) { - log.debug("[{}] Error getting policies while starting repl-producers {}", topic, ex.getMessage()); + List configuredClusters = topicPolicies.getReplicationClusters().get(); + replicators.forEach((region, replicator) -> { + if (configuredClusters.contains(region)) { + replicator.startProducer(); } - replicators.forEach((region, replicator) -> replicator.startProducer()); - return null; }); + return CompletableFuture.completedFuture(null); } public CompletableFuture stopReplProducers() { @@ -3036,8 +3020,7 @@ public void onUpdate(TopicPolicies policies) { subscribeRateLimiter.ifPresent(subscribeRateLimiter -> subscribeRateLimiter.onSubscribeRateUpdate(getSubscribeRate())); } - replicators.forEach((name, replicator) -> replicator.getRateLimiter() - .ifPresent(DispatchRateLimiter::updateDispatchRate)); + checkReplicator(); checkDeduplicationStatus(); @@ -3052,6 +3035,27 @@ public void onUpdate(TopicPolicies policies) { }); } + private CompletableFuture checkReplicator() { + List> futures = new ArrayList<>(); + List configuredClusters = topicPolicies.getReplicationClusters().get(); + String localCluster = brokerService.pulsar().getConfiguration().getClusterName(); + for (String cluster : configuredClusters) { + if (cluster.equals(localCluster)) { + continue; + } + if (!replicators.containsKey(cluster)) { + futures.add(startReplicator(cluster)); + } + } + replicators.forEach((name, replicator) -> { + replicator.getRateLimiter().ifPresent(DispatchRateLimiter::updateDispatchRate); + if (!configuredClusters.contains(name)) { + futures.add(removeReplicator(name)); + } + }); + return FutureUtil.waitForAll(futures); + } + private Optional getNamespacePolicies() { return DispatchRateLimiter.getPolicies(brokerService, topic); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java index 8811dce221aa3..55839a34d80f6 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java @@ -1427,4 +1427,33 @@ public void testReplicatorWithFailedAck() throws Exception { private static final Logger log = LoggerFactory.getLogger(ReplicatorTest.class); + @Test + public void testTopicPoliciesEnabled() throws Exception { + log.info("--- testTopicPoliciesEnabled ---"); + String namespace = "pulsar/ns2"; + admin1.namespaces().createNamespace(namespace); + admin1.namespaces().setNamespaceReplicationClusters(namespace, Sets.newHashSet("r1", "r2")); + final TopicName dest = TopicName.get( + BrokerTestUtil.newUniqueName("persistent://" + namespace + "/testTopicPoliciesEnabled")); + + @Cleanup + MessageProducer producer1 = new MessageProducer(url1, dest); + log.info("--- Starting producer --- " + url1); + + producer1.produce(2); + + PersistentTopic topic = (PersistentTopic) pulsar1.getBrokerService().getTopic(dest.toString(), false) + .getNow(null).get(); + Awaitility.await().untilAsserted(() -> { + assertTrue(topic.getReplicators().containsKey("r2")); + }); + + admin1.topics().setReplicationClusters(dest.toString(), Lists.newArrayList("r1")); + + Awaitility.await().untilAsserted(() -> { + Set replicationClusters = admin1.topics().getReplicationClusters(dest.toString(), false); + assertTrue(replicationClusters != null && replicationClusters.size() == 1); + assertTrue(topic.getReplicators().isEmpty()); + }); + } } From 65bbb349534bcff50e027cb8613df7a6733043ce Mon Sep 17 00:00:00 2001 From: technoboy Date: Sun, 6 Mar 2022 21:23:36 +0800 Subject: [PATCH 02/18] update method name. --- .../org/apache/pulsar/broker/service/ReplicatorTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java index 55839a34d80f6..857505814d523 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java @@ -1428,13 +1428,13 @@ public void testReplicatorWithFailedAck() throws Exception { private static final Logger log = LoggerFactory.getLogger(ReplicatorTest.class); @Test - public void testTopicPoliciesEnabled() throws Exception { - log.info("--- testTopicPoliciesEnabled ---"); + public void testWhenUpdateReplicationCluster() throws Exception { + log.info("--- testWhenUpdateReplicationCluster ---"); String namespace = "pulsar/ns2"; admin1.namespaces().createNamespace(namespace); admin1.namespaces().setNamespaceReplicationClusters(namespace, Sets.newHashSet("r1", "r2")); final TopicName dest = TopicName.get( - BrokerTestUtil.newUniqueName("persistent://" + namespace + "/testTopicPoliciesEnabled")); + BrokerTestUtil.newUniqueName("persistent://" + namespace + "/testWhenUpdateReplicationCluster")); @Cleanup MessageProducer producer1 = new MessageProducer(url1, dest); From 9dc45a0cfc650345a5768e2ea266c5914fc7cb77 Mon Sep 17 00:00:00 2001 From: technoboy Date: Mon, 7 Mar 2022 13:36:43 +0800 Subject: [PATCH 03/18] fix test. --- .../org/apache/pulsar/broker/service/PersistentTopicTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java index 55ededa773415..0b13d57fe3974 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java @@ -1794,7 +1794,7 @@ public void testAtomicReplicationRemoval() throws Exception { // try to start replicator again topic.startReplProducers().join(); // verify: replicator.startProducer is not invoked - verify(replicator, Mockito.times(1)).startProducer(); + verify(replicator, Mockito.times(0)).startProducer(); // step-3 : complete the callback to remove replicator from the list ArgumentCaptor captor = ArgumentCaptor.forClass(DeleteCursorCallback.class); From 99cd941ecdf633ff9fd77b04144aa8336e95ddb9 Mon Sep 17 00:00:00 2001 From: technoboy Date: Mon, 7 Mar 2022 15:29:38 +0800 Subject: [PATCH 04/18] updates. --- .../service/persistent/PersistentTopic.java | 27 ++++--------------- 1 file changed, 5 insertions(+), 22 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 7c292dae12bea..693c7beeebb63 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -3020,7 +3020,11 @@ public void onUpdate(TopicPolicies policies) { subscribeRateLimiter.ifPresent(subscribeRateLimiter -> subscribeRateLimiter.onSubscribeRateUpdate(getSubscribeRate())); } - checkReplicator(); + + checkReplicationAndRetryOnFailure().thenAccept(__ -> { + replicators.forEach((name, replicator) -> replicator.getRateLimiter() + .ifPresent(DispatchRateLimiter::updateDispatchRate)); + }); checkDeduplicationStatus(); @@ -3035,27 +3039,6 @@ public void onUpdate(TopicPolicies policies) { }); } - private CompletableFuture checkReplicator() { - List> futures = new ArrayList<>(); - List configuredClusters = topicPolicies.getReplicationClusters().get(); - String localCluster = brokerService.pulsar().getConfiguration().getClusterName(); - for (String cluster : configuredClusters) { - if (cluster.equals(localCluster)) { - continue; - } - if (!replicators.containsKey(cluster)) { - futures.add(startReplicator(cluster)); - } - } - replicators.forEach((name, replicator) -> { - replicator.getRateLimiter().ifPresent(DispatchRateLimiter::updateDispatchRate); - if (!configuredClusters.contains(name)) { - futures.add(removeReplicator(name)); - } - }); - return FutureUtil.waitForAll(futures); - } - private Optional getNamespacePolicies() { return DispatchRateLimiter.getPolicies(brokerService, topic); } From 5f47395fd58e9654b72f6f3065ae23f76da017ef Mon Sep 17 00:00:00 2001 From: technoboy Date: Mon, 7 Mar 2022 15:34:20 +0800 Subject: [PATCH 05/18] updates. --- .../pulsar/broker/service/persistent/PersistentTopic.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 693c7beeebb63..987bf8d721e4e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -3021,7 +3021,7 @@ public void onUpdate(TopicPolicies policies) { subscribeRateLimiter.onSubscribeRateUpdate(getSubscribeRate())); } - checkReplicationAndRetryOnFailure().thenAccept(__ -> { + checkReplicationAndRetryOnFailure().whenComplete((v, ex) -> { replicators.forEach((name, replicator) -> replicator.getRateLimiter() .ifPresent(DispatchRateLimiter::updateDispatchRate)); }); From e46ccc89cf1fe837d3210a907f45a00ea750e8f9 Mon Sep 17 00:00:00 2001 From: technoboy Date: Tue, 8 Mar 2022 15:08:34 +0800 Subject: [PATCH 06/18] updates. --- .../broker/service/persistent/PersistentTopic.java | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 987bf8d721e4e..dda4dd9c73181 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -3021,10 +3021,12 @@ public void onUpdate(TopicPolicies policies) { subscribeRateLimiter.onSubscribeRateUpdate(getSubscribeRate())); } - checkReplicationAndRetryOnFailure().whenComplete((v, ex) -> { - replicators.forEach((name, replicator) -> replicator.getRateLimiter() - .ifPresent(DispatchRateLimiter::updateDispatchRate)); - }); + replicators.forEach((name, replicator) -> replicator.getRateLimiter() + .ifPresent(DispatchRateLimiter::updateDispatchRate)); + + if (policies.getReplicationClusters() != null) { + checkReplicationAndRetryOnFailure(); + } checkDeduplicationStatus(); From f96e2fc0ae2f4cf41929b67e2e8adf61732bc399 Mon Sep 17 00:00:00 2001 From: technoboy Date: Tue, 8 Mar 2022 15:09:46 +0800 Subject: [PATCH 07/18] updates. --- .../apache/pulsar/broker/service/persistent/PersistentTopic.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index dda4dd9c73181..c21bc5d7fd19c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -3020,7 +3020,6 @@ public void onUpdate(TopicPolicies policies) { subscribeRateLimiter.ifPresent(subscribeRateLimiter -> subscribeRateLimiter.onSubscribeRateUpdate(getSubscribeRate())); } - replicators.forEach((name, replicator) -> replicator.getRateLimiter() .ifPresent(DispatchRateLimiter::updateDispatchRate)); From a002642b61949a3720724156c87b52dd64510809 Mon Sep 17 00:00:00 2001 From: technoboy Date: Fri, 11 Mar 2022 21:16:20 +0800 Subject: [PATCH 08/18] fix test. --- .../java/org/apache/pulsar/broker/service/ReplicatorTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java index 857505814d523..125098809f159 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java @@ -717,7 +717,7 @@ public void testReplicatorProducerClosing() throws Exception { assertNull(producer); } - @Test(priority = 5, timeOut = 30000) + @Test(priority = 4, timeOut = 30000) public void testReplicatorProducerName() throws Exception { log.info("--- Starting ReplicatorTest::testReplicatorProducerName ---"); final String topicName = BrokerTestUtil.newUniqueName("persistent://pulsar/ns/testReplicatorProducerName"); @@ -739,7 +739,7 @@ public void testReplicatorProducerName() throws Exception { }); } - @Test(priority = 5, timeOut = 30000) + @Test(priority = 4, timeOut = 30000) public void testReplicatorProducerNameWithUserDefinedReplicatorPrefix() throws Exception { log.info("--- Starting ReplicatorTest::testReplicatorProducerNameWithUserDefinedReplicatorPrefix ---"); final String topicName = BrokerTestUtil.newUniqueName( From 9ffc5fbb34aeeaa51945630a17cd167ddfb0f33a Mon Sep 17 00:00:00 2001 From: technoboy Date: Sun, 6 Mar 2022 21:18:58 +0800 Subject: [PATCH 09/18] Fix update replication cluster but not update replicator --- .../service/persistent/PersistentTopic.java | 24 +++++++++++++++++-- .../pulsar/broker/service/ReplicatorTest.java | 1 - 2 files changed, 22 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index c21bc5d7fd19c..5a0a76f7423ee 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -3020,8 +3020,7 @@ public void onUpdate(TopicPolicies policies) { subscribeRateLimiter.ifPresent(subscribeRateLimiter -> subscribeRateLimiter.onSubscribeRateUpdate(getSubscribeRate())); } - replicators.forEach((name, replicator) -> replicator.getRateLimiter() - .ifPresent(DispatchRateLimiter::updateDispatchRate)); + checkReplicator(); if (policies.getReplicationClusters() != null) { checkReplicationAndRetryOnFailure(); @@ -3040,6 +3039,27 @@ public void onUpdate(TopicPolicies policies) { }); } + private CompletableFuture checkReplicator() { + List> futures = new ArrayList<>(); + List configuredClusters = topicPolicies.getReplicationClusters().get(); + String localCluster = brokerService.pulsar().getConfiguration().getClusterName(); + for (String cluster : configuredClusters) { + if (cluster.equals(localCluster)) { + continue; + } + if (!replicators.containsKey(cluster)) { + futures.add(startReplicator(cluster)); + } + } + replicators.forEach((name, replicator) -> { + replicator.getRateLimiter().ifPresent(DispatchRateLimiter::updateDispatchRate); + if (!configuredClusters.contains(name)) { + futures.add(removeReplicator(name)); + } + }); + return FutureUtil.waitForAll(futures); + } + private Optional getNamespacePolicies() { return DispatchRateLimiter.getPolicies(brokerService, topic); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java index 125098809f159..c31373c05fcd6 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java @@ -1435,7 +1435,6 @@ public void testWhenUpdateReplicationCluster() throws Exception { admin1.namespaces().setNamespaceReplicationClusters(namespace, Sets.newHashSet("r1", "r2")); final TopicName dest = TopicName.get( BrokerTestUtil.newUniqueName("persistent://" + namespace + "/testWhenUpdateReplicationCluster")); - @Cleanup MessageProducer producer1 = new MessageProducer(url1, dest); log.info("--- Starting producer --- " + url1); From 2b8b172dfa1c3a9bbcb1134718ad013897ea054d Mon Sep 17 00:00:00 2001 From: technoboy Date: Mon, 7 Mar 2022 15:29:38 +0800 Subject: [PATCH 10/18] updates. --- .../service/persistent/PersistentTopic.java | 27 ++++--------------- 1 file changed, 5 insertions(+), 22 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 5a0a76f7423ee..21c67dc371bba 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -3020,7 +3020,11 @@ public void onUpdate(TopicPolicies policies) { subscribeRateLimiter.ifPresent(subscribeRateLimiter -> subscribeRateLimiter.onSubscribeRateUpdate(getSubscribeRate())); } - checkReplicator(); + + checkReplicationAndRetryOnFailure().thenAccept(__ -> { + replicators.forEach((name, replicator) -> replicator.getRateLimiter() + .ifPresent(DispatchRateLimiter::updateDispatchRate)); + }); if (policies.getReplicationClusters() != null) { checkReplicationAndRetryOnFailure(); @@ -3039,27 +3043,6 @@ public void onUpdate(TopicPolicies policies) { }); } - private CompletableFuture checkReplicator() { - List> futures = new ArrayList<>(); - List configuredClusters = topicPolicies.getReplicationClusters().get(); - String localCluster = brokerService.pulsar().getConfiguration().getClusterName(); - for (String cluster : configuredClusters) { - if (cluster.equals(localCluster)) { - continue; - } - if (!replicators.containsKey(cluster)) { - futures.add(startReplicator(cluster)); - } - } - replicators.forEach((name, replicator) -> { - replicator.getRateLimiter().ifPresent(DispatchRateLimiter::updateDispatchRate); - if (!configuredClusters.contains(name)) { - futures.add(removeReplicator(name)); - } - }); - return FutureUtil.waitForAll(futures); - } - private Optional getNamespacePolicies() { return DispatchRateLimiter.getPolicies(brokerService, topic); } From e903b468e4d5ff973762c987f18adbfe319fe6ba Mon Sep 17 00:00:00 2001 From: technoboy Date: Mon, 7 Mar 2022 15:34:20 +0800 Subject: [PATCH 11/18] updates. --- .../pulsar/broker/service/persistent/PersistentTopic.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 21c67dc371bba..54d23b5707a1a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -3021,7 +3021,7 @@ public void onUpdate(TopicPolicies policies) { subscribeRateLimiter.onSubscribeRateUpdate(getSubscribeRate())); } - checkReplicationAndRetryOnFailure().thenAccept(__ -> { + checkReplicationAndRetryOnFailure().whenComplete((v, ex) -> { replicators.forEach((name, replicator) -> replicator.getRateLimiter() .ifPresent(DispatchRateLimiter::updateDispatchRate)); }); From e1cb6b1bb9e674daf9dd872950bb68b4aaf71242 Mon Sep 17 00:00:00 2001 From: technoboy Date: Tue, 8 Mar 2022 15:08:34 +0800 Subject: [PATCH 12/18] updates. --- .../broker/service/persistent/PersistentTopic.java | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 54d23b5707a1a..f8d003f4a0fdc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -3021,10 +3021,12 @@ public void onUpdate(TopicPolicies policies) { subscribeRateLimiter.onSubscribeRateUpdate(getSubscribeRate())); } - checkReplicationAndRetryOnFailure().whenComplete((v, ex) -> { - replicators.forEach((name, replicator) -> replicator.getRateLimiter() - .ifPresent(DispatchRateLimiter::updateDispatchRate)); - }); + replicators.forEach((name, replicator) -> replicator.getRateLimiter() + .ifPresent(DispatchRateLimiter::updateDispatchRate)); + + if (policies.getReplicationClusters() != null) { + checkReplicationAndRetryOnFailure(); + } if (policies.getReplicationClusters() != null) { checkReplicationAndRetryOnFailure(); From 6558d8d16ffbfeeddda3d884ef3070a7fce0ee44 Mon Sep 17 00:00:00 2001 From: technoboy Date: Tue, 8 Mar 2022 15:09:46 +0800 Subject: [PATCH 13/18] updates. --- .../apache/pulsar/broker/service/persistent/PersistentTopic.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index f8d003f4a0fdc..0218b668c79ae 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -3020,7 +3020,6 @@ public void onUpdate(TopicPolicies policies) { subscribeRateLimiter.ifPresent(subscribeRateLimiter -> subscribeRateLimiter.onSubscribeRateUpdate(getSubscribeRate())); } - replicators.forEach((name, replicator) -> replicator.getRateLimiter() .ifPresent(DispatchRateLimiter::updateDispatchRate)); From 414ac34643d9d0f2f42a8e80398cf17f7b933adc Mon Sep 17 00:00:00 2001 From: technoboy Date: Wed, 16 Mar 2022 18:11:37 +0800 Subject: [PATCH 14/18] rollback startReplProducers. --- .../service/persistent/PersistentTopic.java | 30 ++++++++++++++----- 1 file changed, 23 insertions(+), 7 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 0218b668c79ae..823a6cafaf1c1 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -27,6 +27,7 @@ import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.Lists; import com.google.common.collect.Maps; +import com.google.common.collect.Sets; import io.netty.buffer.ByteBuf; import io.netty.util.concurrent.FastThreadLocal; import java.time.Clock; @@ -622,13 +623,28 @@ private boolean hasRemoteProducers() { public CompletableFuture startReplProducers() { // read repl-cluster from policies to avoid restart of replicator which are in process of disconnect and close - List configuredClusters = topicPolicies.getReplicationClusters().get(); - replicators.forEach((region, replicator) -> { - if (configuredClusters.contains(region)) { - replicator.startProducer(); - } - }); - return CompletableFuture.completedFuture(null); + return brokerService.pulsar().getPulsarResources().getNamespaceResources() + .getPoliciesAsync(TopicName.get(topic).getNamespaceObject()) + .thenAccept(optPolicies -> { + if (optPolicies.isPresent()) { + if (optPolicies.get().replication_clusters != null) { + Set configuredClusters = Sets.newTreeSet(optPolicies.get().replication_clusters); + replicators.forEach((region, replicator) -> { + if (configuredClusters.contains(region)) { + replicator.startProducer(); + } + }); + } + } else { + replicators.forEach((region, replicator) -> replicator.startProducer()); + } + }).exceptionally(ex -> { + if (log.isDebugEnabled()) { + log.debug("[{}] Error getting policies while starting repl-producers {}", topic, ex.getMessage()); + } + replicators.forEach((region, replicator) -> replicator.startProducer()); + return null; + }); } public CompletableFuture stopReplProducers() { From 8c4f03f78f55ffee8f4238a12272322319e6bf69 Mon Sep 17 00:00:00 2001 From: technoboy Date: Wed, 16 Mar 2022 18:13:11 +0800 Subject: [PATCH 15/18] fix checkstyle. --- .../broker/service/persistent/PersistentTopic.java | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 823a6cafaf1c1..24753e61bbda7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -639,12 +639,12 @@ public CompletableFuture startReplProducers() { replicators.forEach((region, replicator) -> replicator.startProducer()); } }).exceptionally(ex -> { - if (log.isDebugEnabled()) { - log.debug("[{}] Error getting policies while starting repl-producers {}", topic, ex.getMessage()); - } - replicators.forEach((region, replicator) -> replicator.startProducer()); - return null; - }); + if (log.isDebugEnabled()) { + log.debug("[{}] Error getting policies while starting repl-producers {}", topic, ex.getMessage()); + } + replicators.forEach((region, replicator) -> replicator.startProducer()); + return null; + }); } public CompletableFuture stopReplProducers() { From 3f157aa047cdc4dc1edacc473633bfb5a86716d9 Mon Sep 17 00:00:00 2001 From: technoboy Date: Wed, 16 Mar 2022 18:14:14 +0800 Subject: [PATCH 16/18] fix checkstyle. --- .../broker/service/persistent/PersistentTopic.java | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 24753e61bbda7..6af2f30cc3eed 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -639,12 +639,12 @@ public CompletableFuture startReplProducers() { replicators.forEach((region, replicator) -> replicator.startProducer()); } }).exceptionally(ex -> { - if (log.isDebugEnabled()) { - log.debug("[{}] Error getting policies while starting repl-producers {}", topic, ex.getMessage()); - } - replicators.forEach((region, replicator) -> replicator.startProducer()); - return null; - }); + if (log.isDebugEnabled()) { + log.debug("[{}] Error getting policies while starting repl-producers {}", topic, ex.getMessage()); + } + replicators.forEach((region, replicator) -> replicator.startProducer()); + return null; + }); } public CompletableFuture stopReplProducers() { From 76c103b9b571306817c682baa6869e5126c879ec Mon Sep 17 00:00:00 2001 From: technoboy Date: Thu, 17 Mar 2022 09:56:55 +0800 Subject: [PATCH 17/18] update. --- .../pulsar/broker/service/persistent/PersistentTopic.java | 4 ---- 1 file changed, 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 6af2f30cc3eed..069157b2dbe9d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -3043,10 +3043,6 @@ public void onUpdate(TopicPolicies policies) { checkReplicationAndRetryOnFailure(); } - if (policies.getReplicationClusters() != null) { - checkReplicationAndRetryOnFailure(); - } - checkDeduplicationStatus(); preCreateSubscriptionForCompactionIfNeeded(); From 661ec42adfc725306b01e69ebec03cb8eac75931 Mon Sep 17 00:00:00 2001 From: technoboy Date: Thu, 17 Mar 2022 13:47:21 +0800 Subject: [PATCH 18/18] update. --- .../org/apache/pulsar/broker/service/PersistentTopicTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java index 0b13d57fe3974..55ededa773415 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java @@ -1794,7 +1794,7 @@ public void testAtomicReplicationRemoval() throws Exception { // try to start replicator again topic.startReplProducers().join(); // verify: replicator.startProducer is not invoked - verify(replicator, Mockito.times(0)).startProducer(); + verify(replicator, Mockito.times(1)).startProducer(); // step-3 : complete the callback to remove replicator from the list ArgumentCaptor captor = ArgumentCaptor.forClass(DeleteCursorCallback.class);