From 42db69edaaedf1d95403198998de7f38a217cb82 Mon Sep 17 00:00:00 2001 From: Xiaoyu Hou Date: Mon, 20 Jun 2022 18:12:06 +0800 Subject: [PATCH 1/2] Fix removing replica clusters on topic-policies not take effect --- .../service/persistent/PersistentTopic.java | 4 +-- .../service/ReplicatorTopicPoliciesTest.java | 31 +++++++++++++++++++ 2 files changed, 32 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 d8494dae27d04..8440616b4ef8f 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 @@ -3023,9 +3023,7 @@ public void onUpdate(TopicPolicies policies) { updateSubscribeRateLimiter(); replicators.forEach((name, replicator) -> replicator.updateRateLimiter()); checkMessageExpiry(); - if (policies.getReplicationClusters() != null) { - checkReplicationAndRetryOnFailure(); - } + checkReplicationAndRetryOnFailure(); checkDeduplicationStatus(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTopicPoliciesTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTopicPoliciesTest.java index 2a3ce623562c7..4248c8a5cb16b 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTopicPoliciesTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTopicPoliciesTest.java @@ -23,11 +23,15 @@ import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; import com.google.common.collect.Sets; +import java.util.Arrays; import java.util.Collections; import java.util.HashSet; +import java.util.List; import java.util.Set; import java.util.UUID; +import java.util.stream.Collectors; import org.apache.pulsar.broker.PulsarServerException; +import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.api.PulsarClientException; @@ -731,6 +735,33 @@ public void testReplicatorOffloadPolicies() throws Exception { assertNull(admin3.topicPolicies(true).getOffloadPolicies(persistentTopicName))); } + @Test + public void testRemoveReplicationClusters() throws Exception { + final String namespace = "pulsar/partitionedNs-" + UUID.randomUUID(); + final String persistentTopicName = "persistent://" + namespace + "/topic" + UUID.randomUUID(); + + init(namespace, persistentTopicName); + + // set topic replica cluster policy as `r1 r2` + admin1.topics().setReplicationClusters(persistentTopicName, Arrays.asList("r1", "r2")); + PersistentTopic topicRef = + (PersistentTopic) pulsar1.getBrokerService().getTopicReference(persistentTopicName + "-partition-0").get(); + assertNotNull(topicRef); + + Awaitility.await().untilAsserted(() -> { + List replicaClusters = topicRef.getReplicators().keys().stream().sorted().collect(Collectors.toList()); + assertEquals(replicaClusters.size(), 1); + assertEquals(replicaClusters.toString(), "[r2]"); + }); + + // removing topic replica cluster policy, so namespace policy should take effect + admin1.topics().removeReplicationClusters(persistentTopicName); + Awaitility.await().untilAsserted(() -> { + List replicaClusters = topicRef.getReplicators().keys().stream().sorted().collect(Collectors.toList()); + assertEquals(replicaClusters.size(), 2); + assertEquals(replicaClusters.toString(), "[r2, r3]"); + }); + } private void init(String namespace, String topic) throws PulsarAdminException, PulsarClientException, PulsarServerException { From afed43d95d73e300c333ff0788049fe3fc52ed95 Mon Sep 17 00:00:00 2001 From: Xiaoyu Hou Date: Tue, 21 Jun 2022 10:16:23 +0800 Subject: [PATCH 2/2] Fix UT TopicPoliciesTest#testGetSetSubscribeRate --- .../pulsar/broker/service/persistent/PersistentTopic.java | 5 +++-- 1 file changed, 3 insertions(+), 2 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 8440616b4ef8f..741d16b990cdd 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 @@ -305,7 +305,6 @@ public PersistentTopic(String topic, ManagedLedger ledger, BrokerService brokerS @Override public CompletableFuture initialize() { List> futures = new ArrayList<>(); - futures.add(initTopicPolicy()); for (ManagedCursor cursor : ledger.getCursors()) { if (cursor.getName().startsWith(replicatorPrefix)) { String localCluster = brokerService.pulsar().getConfiguration().getClusterName(); @@ -341,7 +340,9 @@ public CompletableFuture initialize() { this.isEncryptionRequired = policies.encryption_required; isAllowAutoUpdateSchema = policies.is_allow_auto_update_schema; - }).exceptionally(ex -> { + }) + .thenCompose(ignore -> initTopicPolicy()) + .exceptionally(ex -> { log.warn("[{}] Error getting policies {} and isEncryptionRequired will be set to false", topic, ex.getMessage()); isEncryptionRequired = false;