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..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; @@ -3023,9 +3024,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 {