From 56b34d69ce2f72e7b0e93f1db9908ba9e5ab8b80 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Wed, 7 Dec 2022 15:53:07 +0100 Subject: [PATCH 1/6] [fix][broker] Fix namespace deletion if __change_events has not been created yet --- .../SystemTopicBasedTopicPoliciesService.java | 16 ++++++++++++++++ .../pulsar/broker/admin/AdminApi2Test.java | 6 +++++- 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java index a24b1eeee101a..69065f6760778 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java @@ -24,6 +24,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList; @@ -51,8 +52,10 @@ import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.policies.data.TopicPolicies; import org.apache.pulsar.common.util.FutureUtil; +import org.apache.pulsar.metadata.api.MetadataStoreException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -106,6 +109,19 @@ private CompletableFuture sendTopicPolicyEvent(TopicName topicName, Action new BrokerServiceException.NotAllowedException("Not allowed to send event to health check topic")); } CompletableFuture result = new CompletableFuture<>(); + try { + final Optional namespacePolicies = pulsarService.getPulsarResources().getNamespaceResources() + .getPolicies(topicName.getNamespaceObject()); + if (namespacePolicies.isPresent() && namespacePolicies.get().deleted) { + log.debug("[{}] skip sending topic policy event since the namespace is deleted", topicName); + result.complete(null); + return result; + } + } catch (MetadataStoreException e) { + result.completeExceptionally(e); + return result; + } + try { createSystemTopicFactoryIfNeeded(); } catch (PulsarServerException e) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java index 7bb35cc4d63a8..64bbca0d3ed1d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java @@ -1648,6 +1648,11 @@ public void testDeleteNamespaceWithTopicPolicies() throws Exception { // create namespace2 String namespace = tenant + "/test-ns2"; + admin.namespaces().createNamespace(namespace, Set.of("test")); + admin.topics().createNonPartitionedTopic(namespace + "/tobedeleted"); + // verify namespace can be deleted even without topic policy events + admin.namespaces().deleteNamespace(namespace, true); + admin.namespaces().createNamespace(namespace, Set.of("test")); // create topic String topic = namespace + "/test-topic2"; @@ -1872,7 +1877,6 @@ public void testListOfNamespaceBundles() throws Exception { @Test public void testForceDeleteNamespace() throws Exception { - conf.setForceDeleteNamespaceAllowed(true); final String namespaceName = "prop-xyz2/ns1"; TenantInfoImpl tenantInfo = new TenantInfoImpl(Set.of("role1", "role2"), Set.of("test")); admin.tenants().createTenant("prop-xyz2", tenantInfo); From 68a7d23b78934e47e0a8d4dce52309a9e3bfc8d7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Fri, 9 Dec 2022 10:22:23 +0100 Subject: [PATCH 2/6] make it async --- .../SystemTopicBasedTopicPoliciesService.java | 88 ++++++++----------- ...temTopicBasedTopicPoliciesServiceTest.java | 11 +++ 2 files changed, 50 insertions(+), 49 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java index 69065f6760778..616360b0f7aa6 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java @@ -24,12 +24,12 @@ import java.util.List; import java.util.Map; import java.util.Objects; -import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Function; import javax.annotation.Nonnull; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.PulsarService; @@ -52,10 +52,8 @@ import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicName; -import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.policies.data.TopicPolicies; import org.apache.pulsar.common.util.FutureUtil; -import org.apache.pulsar.metadata.api.MetadataStoreException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -109,49 +107,44 @@ private CompletableFuture sendTopicPolicyEvent(TopicName topicName, Action new BrokerServiceException.NotAllowedException("Not allowed to send event to health check topic")); } CompletableFuture result = new CompletableFuture<>(); - try { - final Optional namespacePolicies = pulsarService.getPulsarResources().getNamespaceResources() - .getPolicies(topicName.getNamespaceObject()); - if (namespacePolicies.isPresent() && namespacePolicies.get().deleted) { - log.debug("[{}] skip sending topic policy event since the namespace is deleted", topicName); - result.complete(null); - return result; - } - } catch (MetadataStoreException e) { - result.completeExceptionally(e); - return result; - } - - try { - createSystemTopicFactoryIfNeeded(); - } catch (PulsarServerException e) { - result.completeExceptionally(e); - return result; - } - - SystemTopicClient systemTopicClient = - namespaceEventsSystemTopicFactory.createTopicPoliciesSystemTopicClient(topicName.getNamespaceObject()); - - CompletableFuture> writerFuture = systemTopicClient.newWriterAsync(); - writerFuture.whenComplete((writer, ex) -> { - if (ex != null) { - result.completeExceptionally(ex); - } else { - PulsarEvent event = getPulsarEvent(topicName, actionType, policies); - CompletableFuture actionFuture = - ActionType.DELETE.equals(actionType) ? writer.deleteAsync(getEventKey(event), event) - : writer.writeAsync(getEventKey(event), event); - actionFuture.whenComplete(((messageId, e) -> { - if (e != null) { - result.completeExceptionally(e); - } else { + return pulsarService.getPulsarResources().getNamespaceResources() + .getPoliciesAsync(topicName.getNamespaceObject()) + .thenCompose(namespacePolicies -> { + if (namespacePolicies.isPresent() && namespacePolicies.get().deleted) { + log.debug("[{}] skip sending topic policy event since the namespace is deleted", topicName); + result.complete(null); + return new CompletableFuture<>(); + } + return CompletableFuture.completedFuture(null); + }) + .thenCompose(ignore -> { + try { + createSystemTopicFactoryIfNeeded(); + } catch (PulsarServerException e) { + return CompletableFuture.failedFuture(e); + } + SystemTopicClient systemTopicClient = + namespaceEventsSystemTopicFactory.createTopicPoliciesSystemTopicClient( + topicName.getNamespaceObject()); + return systemTopicClient.newWriterAsync(); + }).thenCompose((SystemTopicClient.Writer writer) -> { + PulsarEvent event = getPulsarEvent(topicName, actionType, policies); + final CompletableFuture writeFuture = ActionType.DELETE.equals(actionType) + ? writer.deleteAsync(getEventKey(event), event) : + writer.writeAsync(getEventKey(event), event); + return writeFuture + .handle((messageId, e) -> { + if (e != null) { + return CompletableFuture.failedFuture(e); + } if (messageId != null) { - result.complete(null); + return CompletableFuture.completedFuture(null); } else { - result.completeExceptionally(new RuntimeException("Got message id is null.")); + return CompletableFuture.failedFuture( + new RuntimeException("Got message id is null.")); } - } - writer.closeAsync().whenComplete((v, cause) -> { + }) + .thenRun(() -> writer.closeAsync().whenComplete((v, cause) -> { if (cause != null) { log.error("[{}] Close writer error.", topicName, cause); } else { @@ -159,12 +152,9 @@ private CompletableFuture sendTopicPolicyEvent(TopicName topicName, Action log.debug("[{}] Close writer success.", topicName); } } - }); - }) - ); - } - }); - return result; + })); + }) + .applyToEither(result, Function.identity()); } private PulsarEvent getPulsarEvent(TopicName topicName, ActionType actionType, TopicPolicies policies) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java index 7b66b6a6b5152..f9fc717a817af 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java @@ -373,4 +373,15 @@ public void run() { } }); } + + @Test + public void testHandleNamespaceBeingDeleted() throws Exception { + SystemTopicBasedTopicPoliciesService service = (SystemTopicBasedTopicPoliciesService) pulsar.getTopicPoliciesService(); + pulsar.getPulsarResources().getNamespaceResources().setPolicies(NamespaceName.get(NAMESPACE1), + old -> { + old.deleted = true; + return old; + }); + service.deleteTopicPoliciesAsync(TOPIC1).get(); + } } From c37169ee0e1cf2a8b1e0f0b91ddc5bc12c501759 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Fri, 9 Dec 2022 10:25:58 +0100 Subject: [PATCH 3/6] rm double check about ns policies --- .../apache/pulsar/broker/service/BrokerService.java | 12 ++---------- 1 file changed, 2 insertions(+), 10 deletions(-) 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 c5f3d508a56db..bcec8351733a9 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 @@ -3285,16 +3285,8 @@ public CompletableFuture deleteTopicPolicies(TopicName topicName) { if (!pulsarService.getConfig().isTopicLevelPoliciesEnabled()) { return CompletableFuture.completedFuture(null); } - return pulsarService.getPulsarResources().getNamespaceResources() - .getPoliciesAsync(topicName.getNamespaceObject()) - .thenCompose(optPolicies -> { - if (optPolicies.isPresent() && optPolicies.get().deleted) { - // We can return the completed future directly if the namespace is already deleted. - return CompletableFuture.completedFuture(null); - } - TopicName cloneTopicName = TopicName.get(topicName.getPartitionedTopicName()); - return pulsar.getTopicPoliciesService().deleteTopicPoliciesAsync(cloneTopicName); - }); + return pulsar.getTopicPoliciesService() + .deleteTopicPoliciesAsync(TopicName.get(topicName.getPartitionedTopicName())); } private CompletableFuture checkMaxTopicsPerNamespace(TopicName topicName, int numPartitions) { From 7facf372b8c89c149449c810e57915edea245e6b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Mon, 12 Dec 2022 14:40:18 +0100 Subject: [PATCH 4/6] simplify code --- .../SystemTopicBasedTopicPoliciesService.java | 106 +++++++++--------- 1 file changed, 53 insertions(+), 53 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java index 616360b0f7aa6..481213de5173a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java @@ -24,12 +24,12 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; -import java.util.function.Function; import javax.annotation.Nonnull; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.PulsarService; @@ -52,8 +52,10 @@ import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.policies.data.TopicPolicies; import org.apache.pulsar.common.util.FutureUtil; +import org.apache.pulsar.metadata.api.MetadataStoreException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -106,55 +108,53 @@ private CompletableFuture sendTopicPolicyEvent(TopicName topicName, Action return CompletableFuture.failedFuture( new BrokerServiceException.NotAllowedException("Not allowed to send event to health check topic")); } - CompletableFuture result = new CompletableFuture<>(); return pulsarService.getPulsarResources().getNamespaceResources() .getPoliciesAsync(topicName.getNamespaceObject()) .thenCompose(namespacePolicies -> { if (namespacePolicies.isPresent() && namespacePolicies.get().deleted) { log.debug("[{}] skip sending topic policy event since the namespace is deleted", topicName); - result.complete(null); - return new CompletableFuture<>(); + return CompletableFuture.completedFuture(null); } - return CompletableFuture.completedFuture(null); - }) - .thenCompose(ignore -> { + try { createSystemTopicFactoryIfNeeded(); } catch (PulsarServerException e) { return CompletableFuture.failedFuture(e); } + SystemTopicClient systemTopicClient = - namespaceEventsSystemTopicFactory.createTopicPoliciesSystemTopicClient( - topicName.getNamespaceObject()); - return systemTopicClient.newWriterAsync(); - }).thenCompose((SystemTopicClient.Writer writer) -> { - PulsarEvent event = getPulsarEvent(topicName, actionType, policies); - final CompletableFuture writeFuture = ActionType.DELETE.equals(actionType) - ? writer.deleteAsync(getEventKey(event), event) : - writer.writeAsync(getEventKey(event), event); - return writeFuture - .handle((messageId, e) -> { + namespaceEventsSystemTopicFactory.createTopicPoliciesSystemTopicClient(topicName.getNamespaceObject()); + + return systemTopicClient.newWriterAsync() + .thenCompose(writer -> { + PulsarEvent event = getPulsarEvent(topicName, actionType, policies); + CompletableFuture writeFuture = + ActionType.DELETE.equals(actionType) ? writer.deleteAsync(getEventKey(event), event) + : writer.writeAsync(getEventKey(event), event); + return writeFuture.handle((messageId, e) -> { if (e != null) { return CompletableFuture.failedFuture(e); - } - if (messageId != null) { - return CompletableFuture.completedFuture(null); } else { - return CompletableFuture.failedFuture( - new RuntimeException("Got message id is null.")); - } - }) - .thenRun(() -> writer.closeAsync().whenComplete((v, cause) -> { - if (cause != null) { - log.error("[{}] Close writer error.", topicName, cause); - } else { - if (log.isDebugEnabled()) { - log.debug("[{}] Close writer success.", topicName); + if (messageId != null) { + return CompletableFuture.completedFuture(null); + } else { + return CompletableFuture.failedFuture( + new RuntimeException("Got message id is null.")); } } - })); - }) - .applyToEither(result, Function.identity()); + }).thenRun(() -> + writer.closeAsync().whenComplete((v, cause) -> { + if (cause != null) { + log.error("[{}] Close writer error.", topicName, cause); + } else { + if (log.isDebugEnabled()) { + log.debug("[{}] Close writer success.", topicName); + } + } + }) + ); + }); + }); } private PulsarEvent getPulsarEvent(TopicName topicName, ActionType actionType, TopicPolicies policies) { @@ -396,25 +396,25 @@ private void cleanCacheAndCloseReader(@Nonnull NamespaceName namespace, boolean private void readMorePolicies(SystemTopicClient.Reader reader) { reader.readNextAsync() - .thenAccept(msg -> { - refreshTopicPoliciesCache(msg); - notifyListener(msg); - }) - .whenComplete((__, ex) -> { - if (ex == null) { - readMorePolicies(reader); - } else { - Throwable cause = FutureUtil.unwrapCompletionException(ex); - if (cause instanceof PulsarClientException.AlreadyClosedException) { - log.warn("Read more topic policies exception, close the read now!", ex); - cleanCacheAndCloseReader( - reader.getSystemTopic().getTopicName().getNamespaceObject(), false); - } else { - log.warn("Read more topic polices exception, read again.", ex); - readMorePolicies(reader); - } - } - }); + .thenAccept(msg -> { + refreshTopicPoliciesCache(msg); + notifyListener(msg); + }) + .whenComplete((__, ex) -> { + if (ex == null) { + readMorePolicies(reader); + } else { + Throwable cause = FutureUtil.unwrapCompletionException(ex); + if (cause instanceof PulsarClientException.AlreadyClosedException) { + log.warn("Read more topic policies exception, close the read now!", ex); + cleanCacheAndCloseReader( + reader.getSystemTopic().getTopicName().getNamespaceObject(), false); + } else { + log.warn("Read more topic polices exception, read again.", ex); + readMorePolicies(reader); + } + } + }); } private void refreshTopicPoliciesCache(Message msg) { @@ -483,7 +483,7 @@ private boolean hasReplicateTo(Message message) { if (message instanceof MessageImpl) { return ((MessageImpl) message).hasReplicateTo() ? (((MessageImpl) message).getReplicateTo().size() == 1 - ? !((MessageImpl) message).getReplicateTo().contains(clusterName) : true) + ? !((MessageImpl) message).getReplicateTo().contains(clusterName) : true) : false; } if (message instanceof TopicMessageImpl) { From db8cfafaa93262eddf22a02d09cfe93121b27e4a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Mon, 12 Dec 2022 14:41:46 +0100 Subject: [PATCH 5/6] style --- .../service/SystemTopicBasedTopicPoliciesService.java | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java index 481213de5173a..93a2cddad4754 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java @@ -24,7 +24,6 @@ import java.util.List; import java.util.Map; import java.util.Objects; -import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList; @@ -52,10 +51,8 @@ import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicName; -import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.policies.data.TopicPolicies; import org.apache.pulsar.common.util.FutureUtil; -import org.apache.pulsar.metadata.api.MetadataStoreException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -122,8 +119,8 @@ private CompletableFuture sendTopicPolicyEvent(TopicName topicName, Action return CompletableFuture.failedFuture(e); } - SystemTopicClient systemTopicClient = - namespaceEventsSystemTopicFactory.createTopicPoliciesSystemTopicClient(topicName.getNamespaceObject()); + SystemTopicClient systemTopicClient = namespaceEventsSystemTopicFactory + .createTopicPoliciesSystemTopicClient(topicName.getNamespaceObject()); return systemTopicClient.newWriterAsync() .thenCompose(writer -> { From ac72e176af46865c08bc4ca9ec5374b57a6b0127 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Tue, 13 Dec 2022 15:07:39 +0100 Subject: [PATCH 6/6] fix flaky test --- .../pulsar/broker/admin/AdminApiTest.java | 42 +++++++++++++------ 1 file changed, 29 insertions(+), 13 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java index b303386c65ffd..fec0ce2da6ecf 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java @@ -167,6 +167,10 @@ public class AdminApiTest extends MockedPulsarServiceBaseTest { @BeforeClass @Override public void setup() throws Exception { + setupConfigAndStart(null); + } + + private void applyDefaultConfig() { conf.setSystemTopicEnabled(false); conf.setTopicLevelPoliciesEnabled(false); conf.setLoadBalancerEnabled(true); @@ -178,6 +182,13 @@ public void setup() throws Exception { conf.setSubscriptionExpiryCheckIntervalInMinutes(1); conf.setBrokerDeleteInactiveTopicsEnabled(false); conf.setNumExecutorThreadPoolSize(5); + } + + private void setupConfigAndStart(java.util.function.Consumer configurationConsumer) throws Exception { + applyDefaultConfig(); + if (configurationConsumer != null) { + configurationConsumer.accept(conf); + } super.internalSetup(); @@ -215,6 +226,7 @@ public void resetClusters() throws Exception { pulsar.getConfiguration().setForceDeleteNamespaceAllowed(false); resetConfig(); + applyDefaultConfig(); setupClusters(); } @@ -1718,9 +1730,9 @@ public void testNamespaceSplitBundleWithInvalidAlgorithm() { @Test public void testNamespaceSplitBundleWithDefaultTopicCountEquallyDivideAlgorithm() throws Exception { cleanup(); - setup(); + setupConfigAndStart(conf -> conf + .setDefaultNamespaceBundleSplitAlgorithm(NamespaceBundleSplitAlgorithm.TOPIC_COUNT_EQUALLY_DIVIDE)); - conf.setDefaultNamespaceBundleSplitAlgorithm(NamespaceBundleSplitAlgorithm.TOPIC_COUNT_EQUALLY_DIVIDE); // Force to create a topic final String namespace = "prop-xyz/ns1"; List topicNames = Lists.newArrayList( @@ -1756,7 +1768,6 @@ public void testNamespaceSplitBundleWithDefaultTopicCountEquallyDivideAlgorithm( for (Producer producer : producers) { producer.close(); } - conf.setDefaultNamespaceBundleSplitAlgorithm(NamespaceBundleSplitAlgorithm.RANGE_EQUALLY_DIVIDE_NAME); } @Test @@ -1911,9 +1922,6 @@ public void testNamespaceUnloadBundle() throws Exception { @Test(dataProvider = "numBundles") public void testNamespaceBundleUnload(Integer numBundles) throws Exception { - cleanup(); - setup(); - admin.namespaces().createNamespace("prop-xyz/ns1-bundles", numBundles); admin.namespaces().setNamespaceReplicationClusters("prop-xyz/ns1-bundles", Set.of("test")); @@ -3259,7 +3267,8 @@ public void testBacklogSizeShouldBeZeroWhenConsumerAckedAllMessages() throws Exc @Test public void testGetTtlDurationDefaultInSeconds() throws Exception { - conf.setTtlDurationDefaultInSeconds(3600); + cleanup(); + setupConfigAndStart(conf -> conf.setTtlDurationDefaultInSeconds(3600)); Integer seconds = admin.namespaces().getPolicies("prop-xyz/ns1").message_ttl_in_seconds; assertNull(seconds); } @@ -3309,8 +3318,11 @@ public void testPartitionedTopicMsgDelayedAggregated() throws Exception { final String topic = "persistent://prop-xyz/ns1/testPartitionedTopicMsgDelayedAggregated-" + UUID.randomUUID().toString(); final String subName = "my-sub"; final int numPartitions = 2; - conf.setSubscriptionRedeliveryTrackerEnabled(true); - conf.setDelayedDeliveryEnabled(true); + cleanup(); + setupConfigAndStart(conf -> { + conf.setSubscriptionRedeliveryTrackerEnabled(true); + conf.setDelayedDeliveryEnabled(true); + }); admin.topics().createPartitionedTopic(topic, numPartitions); for (int i = 0; i < 2; i++) { @@ -3367,7 +3379,7 @@ public void testPartitionedTopicMsgDelayedAggregated() throws Exception { @Test(timeOut = 20000) public void testPartitionedTopicTruncate() throws Exception { - final String topicName = "persistent://prop-xyz/ns1/testTruncateTopic-" + UUID.randomUUID().toString(); + final String topicName = "persistent://prop-xyz/ns1/testTruncateTopic2-" + UUID.randomUUID().toString(); final String subName = "my-sub"; admin.topics().createPartitionedTopic(topicName,6); admin.namespaces().setRetention("prop-xyz/ns1", new RetentionPolicies(60, 50)); @@ -3387,9 +3399,13 @@ public void testPartitionedTopicTruncate() throws Exception { @Test(timeOut = 20000) public void testNonPartitionedTopicTruncate() throws Exception { - final String topicName = "persistent://prop-xyz/ns1/testTruncateTopic-" + UUID.randomUUID().toString(); + final String topicName = "persistent://prop-xyz/ns1/testTruncateTopic1-" + UUID.randomUUID().toString(); final String subName = "my-sub"; - this.conf.setTopicLevelPoliciesEnabled(true); + cleanup(); + setupConfigAndStart(conf -> { + conf.setTopicLevelPoliciesEnabled(true); + conf.setSystemTopicEnabled(true); + }); admin.topics().createNonPartitionedTopic(topicName); admin.namespaces().setRetention("prop-xyz/ns1", new RetentionPolicies(60, 50)); List messageIds = publishMessagesOnPersistentTopic(topicName, 10); @@ -3405,7 +3421,7 @@ public void testNonPartitionedTopicTruncate() throws Exception { @Test(timeOut = 20000) public void testNonPersistentTopicTruncate() throws Exception { - final String topicName = "non-persistent://prop-xyz/ns1/testTruncateTopic-" + UUID.randomUUID().toString(); + final String topicName = "non-persistent://prop-xyz/ns1/testTruncateTopic3-" + UUID.randomUUID().toString(); admin.topics().createNonPartitionedTopic(topicName); assertThrows(() -> {admin.topics().truncate(topicName);}); }