From 3eee3613be235781e66a4c0587dd616efaa797f0 Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Wed, 2 Nov 2022 22:54:06 +0800 Subject: [PATCH 1/6] [fix][broker] Fix can not delete namespace by force ### Motivation 1. When we delete a namespace by force, if the __transaction_buffer_snapshot and __change_events be deleted first, and then the deleting of the normal topic will be failed at clearing snapshot and clearing topic policies. 2. fix flaky test testCreateTransactionSystemTopic. Test failed due to cleaning up after class. ### Modification 1. Delete normal topic and then delete system topic. 2. The system topic is not need to clean up topic policies. --- .../broker/admin/impl/NamespacesBase.java | 43 +++++++++++-------- .../service/persistent/PersistentTopic.java | 9 +++- .../broker/transaction/TransactionTest.java | 5 ++- 3 files changed, 37 insertions(+), 20 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java index 4dcd8809c7788..ad2550e0d4597 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java @@ -218,24 +218,29 @@ protected CompletableFuture internalDeleteNamespaceAsync(boolean force) { })) .thenCompose(topics -> { List allTopics = topics.get(0); + ArrayList copyAllTopics = new ArrayList<>(); List allPartitionedTopics = topics.get(1); - if (!force) { - boolean hasNonSystemTopic = false; - for (String topic : allTopics) { - if (!pulsar().getBrokerService().isSystemTopic(TopicName.get(topic))) { - hasNonSystemTopic = true; - break; - } + ArrayList copyAllPartitionTopics = new ArrayList<>(); + boolean hasNonSystemTopic = false; + List allSystemTopics = new ArrayList<>(); + List allPartitionedSystemTopics = new ArrayList<>(); + for (String topic : allTopics) { + if (!pulsar().getBrokerService().isSystemTopic(TopicName.get(topic))) { + hasNonSystemTopic = true; + copyAllTopics.add(topic); + } else { + allSystemTopics.add(topic); } - if (!hasNonSystemTopic) { - for (String topic : allPartitionedTopics) { - if (!pulsar().getBrokerService().isSystemTopic(TopicName.get(topic))) { - hasNonSystemTopic = true; - break; - } - } + } + for (String topic : allPartitionedTopics) { + if (!pulsar().getBrokerService().isSystemTopic(TopicName.get(topic))) { + hasNonSystemTopic = true; + copyAllPartitionTopics.add(topic); + } else { + allPartitionedSystemTopics.add(topic); } - + } + if (!force) { if (hasNonSystemTopic) { throw new RestException(Status.CONFLICT, "Cannot delete non empty namespace"); } @@ -244,9 +249,13 @@ protected CompletableFuture internalDeleteNamespaceAsync(boolean force) { old.deleted = true; return old; }).thenCompose(__ -> { - return internalDeleteTopicsAsync(allTopics); + return internalDeleteTopicsAsync(copyAllTopics); + }).thenCompose(__ -> { + return internalDeletePartitionedTopicsAsync(copyAllPartitionTopics); + }).thenCompose(__ -> { + return internalDeleteTopicsAsync(allSystemTopics); }).thenCompose(__ -> { - return internalDeletePartitionedTopicsAsync(allPartitionedTopics); + return internalDeletePartitionedTopicsAsync(allPartitionedSystemTopics); }); }) .thenCompose(__ -> pulsar().getNamespaceService() 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 965c1e164a3fa..03b0c1db409dc 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 @@ -1207,7 +1207,14 @@ private CompletableFuture delete(boolean failIfHasSubscriptions, brokerService.deleteTopicAuthenticationWithRetry(topic, deleteTopicAuthenticationFuture, 5); deleteTopicAuthenticationFuture.thenCompose(__ -> deleteSchema()) - .thenCompose(__ -> deleteTopicPolicies()) + .thenCompose(__ -> { + if (!this.getBrokerService().getPulsar().getBrokerService() + .isSystemTopic(TopicName.get(topic))) { + return deleteTopicPolicies(); + } else { + return CompletableFuture.completedFuture(null); + } + }) .thenCompose(__ -> transactionBufferCleanupAndClose()) .whenComplete((v, ex) -> { if (ex != null) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index 33305a2b8df10..6a47f3e3389b4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -215,7 +215,8 @@ public void testTopicTransactionMetrics() throws Exception { public void testCreateTransactionSystemTopic() throws Exception { String subName = "test"; String topicName = TopicName.get(NAMESPACE1 + "/" + "testCreateTransactionSystemTopic").toString(); - + admin.namespaces().deleteNamespace(NAMESPACE1, true); + admin.namespaces().createNamespace(NAMESPACE1); try { // init pending ack @Cleanup @@ -231,7 +232,7 @@ public void testCreateTransactionSystemTopic() throws Exception { // getList does not include transaction system topic List list = admin.topics().getList(NAMESPACE1); - assertEquals(list.size(), 2); + assertEquals(list.size(), 1); list.forEach(topic -> assertFalse(topic.contains(PENDING_ACK_STORE_SUFFIX))); try { From 82610fd7eda71e0f7da9fa1566437c1a8dd20722 Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Thu, 3 Nov 2022 15:12:05 +0800 Subject: [PATCH 2/6] fix some comments, and delete a needless test. --- .../transaction/TransactionProduceTest.java | 29 ------------------- .../broker/transaction/TransactionTest.java | 12 +++++++- 2 files changed, 11 insertions(+), 30 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionProduceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionProduceTest.java index 06bb92890b6c6..cdbb1563280b4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionProduceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionProduceTest.java @@ -89,35 +89,6 @@ public void produceAndCommitTest() throws Exception { produceTest(true); } - @Test - public void testDeleteNamespaceBeforeCommit() throws Exception { - final String topic = NAMESPACE1 + "/testDeleteTopicBeforeCommit"; - PulsarClient pulsarClient = this.pulsarClient; - Transaction tnx = pulsarClient.newTransaction() - .withTransactionTimeout(60, TimeUnit.SECONDS) - .build().get(); - long txnIdMostBits = ((TransactionImpl) tnx).getTxnIdMostBits(); - long txnIdLeastBits = ((TransactionImpl) tnx).getTxnIdLeastBits(); - Assert.assertTrue(txnIdMostBits > -1); - Assert.assertTrue(txnIdLeastBits > -1); - - @Cleanup - Producer outProducer = pulsarClient - .newProducer() - .topic(topic) - .sendTimeout(0, TimeUnit.SECONDS) - .enableBatching(false) - .create(); - - String content = "Hello Txn"; - outProducer.newMessage(tnx).value(content.getBytes(UTF_8)).send(); - - try { - deleteNamespaceGraceFully(NAMESPACE1, true); - } catch (Exception ignore) {} - tnx.commit().get(); - } - @Test public void produceAndAbortTest() throws Exception { produceTest(false); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index 6a47f3e3389b4..1d48f2559bef1 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -211,6 +211,16 @@ public void testTopicTransactionMetrics() throws Exception { assertEquals(stats.ongoingTxnCount, 1); } + @Test + public void testDeleteNamespaceAfterUsedTransaction() throws Exception { + String topicName = TopicName.get(NAMESPACE1 + "/" + "testDeleteNamespaceAfterUsedTransaction").toString(); + Producer producer = pulsarClient.newProducer() + .topic(topicName) + .create(); + + admin.namespaces().deleteNamespace(NAMESPACE1, true); + } + @Test public void testCreateTransactionSystemTopic() throws Exception { String subName = "test"; @@ -232,7 +242,7 @@ public void testCreateTransactionSystemTopic() throws Exception { // getList does not include transaction system topic List list = admin.topics().getList(NAMESPACE1); - assertEquals(list.size(), 1); + assertFalse(list.isEmpty()); list.forEach(topic -> assertFalse(topic.contains(PENDING_ACK_STORE_SUFFIX))); try { From dffda40de622907a2db271b6d79328c49a9207cd Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Thu, 3 Nov 2022 20:56:32 +0800 Subject: [PATCH 3/6] optimize --- .../pulsar/broker/admin/impl/NamespacesBase.java | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java index ad2550e0d4597..6da70e234fbef 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java @@ -248,17 +248,17 @@ protected CompletableFuture internalDeleteNamespaceAsync(boolean force) { return namespaceResources().setPoliciesAsync(namespaceName, old -> { old.deleted = true; return old; - }).thenCompose(__ -> { + }).thenCompose(ignore -> { return internalDeleteTopicsAsync(copyAllTopics); - }).thenCompose(__ -> { + }).thenCompose(ignore -> { return internalDeletePartitionedTopicsAsync(copyAllPartitionTopics); - }).thenCompose(__ -> { + }).thenCompose(ignore -> { return internalDeleteTopicsAsync(allSystemTopics); - }).thenCompose(__ -> { + }).thenCompose(ignore__ -> { return internalDeletePartitionedTopicsAsync(allPartitionedSystemTopics); }); }) - .thenCompose(__ -> pulsar().getNamespaceService() + .thenCompose(ignore -> pulsar().getNamespaceService() .getNamespaceBundleFactory().getBundlesAsync(namespaceName)) .thenCompose(bundles -> FutureUtil.waitForAll(bundles.getBundles().stream() .map(bundle -> pulsar().getNamespaceService().getOwnerAsync(bundle) @@ -280,7 +280,7 @@ protected CompletableFuture internalDeleteNamespaceAsync(boolean force) { return CompletableFuture.completedFuture(null); }) ).collect(Collectors.toList()))) - .thenCompose(__ -> internalClearZkSources()); + .thenCompose(ignore -> internalClearZkSources()); } private CompletableFuture internalDeletePartitionedTopicsAsync(List topicNames) { From 9b919478fe4619d52437ed6214bb4b3488f0838a Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Fri, 4 Nov 2022 10:32:56 +0800 Subject: [PATCH 4/6] fix some comments --- .../pulsar/broker/admin/impl/NamespacesBase.java | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java index 6da70e234fbef..eb2b9d6367ac2 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java @@ -218,16 +218,16 @@ protected CompletableFuture internalDeleteNamespaceAsync(boolean force) { })) .thenCompose(topics -> { List allTopics = topics.get(0); - ArrayList copyAllTopics = new ArrayList<>(); + ArrayList allUserCreatedTopics = new ArrayList<>(); List allPartitionedTopics = topics.get(1); - ArrayList copyAllPartitionTopics = new ArrayList<>(); + ArrayList allUserCreatedPartitionTopics = new ArrayList<>(); boolean hasNonSystemTopic = false; List allSystemTopics = new ArrayList<>(); List allPartitionedSystemTopics = new ArrayList<>(); for (String topic : allTopics) { if (!pulsar().getBrokerService().isSystemTopic(TopicName.get(topic))) { hasNonSystemTopic = true; - copyAllTopics.add(topic); + allUserCreatedTopics.add(topic); } else { allSystemTopics.add(topic); } @@ -235,7 +235,7 @@ protected CompletableFuture internalDeleteNamespaceAsync(boolean force) { for (String topic : allPartitionedTopics) { if (!pulsar().getBrokerService().isSystemTopic(TopicName.get(topic))) { hasNonSystemTopic = true; - copyAllPartitionTopics.add(topic); + allUserCreatedPartitionTopics.add(topic); } else { allPartitionedSystemTopics.add(topic); } @@ -249,9 +249,9 @@ protected CompletableFuture internalDeleteNamespaceAsync(boolean force) { old.deleted = true; return old; }).thenCompose(ignore -> { - return internalDeleteTopicsAsync(copyAllTopics); + return internalDeleteTopicsAsync(allUserCreatedTopics); }).thenCompose(ignore -> { - return internalDeletePartitionedTopicsAsync(copyAllPartitionTopics); + return internalDeletePartitionedTopicsAsync(allUserCreatedPartitionTopics); }).thenCompose(ignore -> { return internalDeleteTopicsAsync(allSystemTopics); }).thenCompose(ignore__ -> { From 68155565df449736fae24a07f52dacc1f3779164 Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Fri, 4 Nov 2022 12:36:53 +0800 Subject: [PATCH 5/6] fix some comments --- .../pulsar/broker/admin/NamespacesTest.java | 64 +++++++++++++++++++ .../broker/transaction/TransactionTest.java | 10 --- 2 files changed, 64 insertions(+), 10 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java index 307c8447674f6..9742e3361f140 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java @@ -76,6 +76,7 @@ import org.apache.pulsar.broker.namespace.NamespaceService; import org.apache.pulsar.broker.namespace.OwnershipCache; import org.apache.pulsar.broker.service.AbstractTopic; +import org.apache.pulsar.broker.service.SystemTopicBasedTopicPoliciesService; import org.apache.pulsar.broker.service.Topic; import org.apache.pulsar.broker.web.PulsarWebResource; import org.apache.pulsar.broker.web.RestException; @@ -92,6 +93,8 @@ import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.NamespaceBundles; import org.apache.pulsar.common.naming.NamespaceName; +import org.apache.pulsar.common.naming.SystemTopicNames; +import org.apache.pulsar.common.naming.TopicDomain; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.AutoTopicCreationOverride; @@ -119,6 +122,7 @@ import org.mockito.Mockito; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.testng.Assert; import org.testng.annotations.AfterClass; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeClass; @@ -1998,4 +2002,64 @@ private void cleanupNamespaceByNsCollection(Collection namespaces) pulsar.getConfiguration().setForceDeleteNamespaceAllowed(forceDeleteNamespaceAllowedOriginalValue); } + @Test + public void testFinallyDeleteSystemTopicWhenDeleteNamespace() throws Exception { + String namespace = this.testTenant + "/delete-namespace"; + String topic = TopicName.get(TopicDomain.persistent.toString(), this.testTenant, "delete-namespace", + "testFinallyDeleteSystemTopicWhenDeleteNamespace").toString(); + + // 0. enable topic level polices and system topic + pulsar.getConfig().setTopicLevelPoliciesEnabled(true); + pulsar.getConfig().setSystemTopicEnabled(true); + pulsar.getConfig().setForceDeleteNamespaceAllowed(true); + Field policesService = pulsar.getClass().getDeclaredField("topicPoliciesService"); + policesService.setAccessible(true); + policesService.set(pulsar, new SystemTopicBasedTopicPoliciesService(pulsar)); + + // 1. create a test namespace. + admin.namespaces().createNamespace(namespace); + // 2. create a test topic. + admin.topics().createNonPartitionedTopic(topic); + // 3. change policy of the topic. + admin.topicPolicies().setMaxConsumers(topic, 5); + // 4. change the order of the topics in this namespace. + List topics = pulsar.getNamespaceService().getFullListOfTopics(NamespaceName.get(namespace)).get(); + Assert.assertTrue(topics.size() >= 2); + for (int i = 0; i < topics.size(); i++) { + if (topics.get(i).contains(SystemTopicNames.NAMESPACE_EVENTS_LOCAL_NAME)) { + String systemTopic = topics.get(i); + topics.set(i, topics.get(0)); + topics.set(0, systemTopic); + } + } + NamespaceService mockNamespaceService = spy(pulsar.getNamespaceService()); + Field namespaceServiceField = pulsar.getClass().getDeclaredField("nsService"); + namespaceServiceField.setAccessible(true); + namespaceServiceField.set(pulsar, mockNamespaceService); + doReturn(CompletableFuture.completedFuture(topics)).when(mockNamespaceService).getFullListOfTopics(any()); + // 5. delete the namespace + admin.namespaces().deleteNamespace(namespace, true); + } + + @Test + public void testNotClearTopicPolicesWhenDeleteSystemTopic() throws Exception { + String namespace = this.testTenant + "/delete-systemTopic"; + String topic = TopicName.get(TopicDomain.persistent.toString(), this.testTenant, "delete-systemTopic", + "testNotClearTopicPolicesWhenDeleteSystemTopic").toString(); + + // 0. enable topic level polices and system topic + pulsar.getConfig().setTopicLevelPoliciesEnabled(true); + pulsar.getConfig().setSystemTopicEnabled(true); + Field policesService = pulsar.getClass().getDeclaredField("topicPoliciesService"); + policesService.setAccessible(true); + policesService.set(pulsar, new SystemTopicBasedTopicPoliciesService(pulsar)); + // 1. create a test namespace. + admin.namespaces().createNamespace(namespace); + // 2. create a test topic. + admin.topics().createNonPartitionedTopic(topic); + // 3. change policy of the topic. + admin.topicPolicies().setMaxConsumers(topic, 5); + // 4. delete the policies topic and the topic wil not to clear topic polices + admin.topics().delete(namespace + "/" + SystemTopicNames.NAMESPACE_EVENTS_LOCAL_NAME, true); + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index 1d48f2559bef1..1dd6feb4762ca 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -211,16 +211,6 @@ public void testTopicTransactionMetrics() throws Exception { assertEquals(stats.ongoingTxnCount, 1); } - @Test - public void testDeleteNamespaceAfterUsedTransaction() throws Exception { - String topicName = TopicName.get(NAMESPACE1 + "/" + "testDeleteNamespaceAfterUsedTransaction").toString(); - Producer producer = pulsarClient.newProducer() - .topic(topicName) - .create(); - - admin.namespaces().deleteNamespace(NAMESPACE1, true); - } - @Test public void testCreateTransactionSystemTopic() throws Exception { String subName = "test"; From a5ef7a3b003e9b836da48e3468aa8b036d6c8eb1 Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Sat, 5 Nov 2022 22:41:20 +0800 Subject: [PATCH 6/6] clean up after test --- .../java/org/apache/pulsar/broker/admin/NamespacesTest.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java index 9742e3361f140..caad3bba7d003 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java @@ -2039,6 +2039,8 @@ public void testFinallyDeleteSystemTopicWhenDeleteNamespace() throws Exception { doReturn(CompletableFuture.completedFuture(topics)).when(mockNamespaceService).getFullListOfTopics(any()); // 5. delete the namespace admin.namespaces().deleteNamespace(namespace, true); + // cleanup + resetBroker(); } @Test