From 0556b273e31cfbb8a56769284f01950f219c0a01 Mon Sep 17 00:00:00 2001 From: congbo Date: Tue, 22 Jun 2021 23:40:48 +0800 Subject: [PATCH 1/6] [Transaction] Fix delete sub then delete pending ack. --- .../service/persistent/PersistentTopic.java | 93 +++++++++++++------ .../pendingack/PendingAckPersistentTest.java | 70 +++++++++++++- 2 files changed, 134 insertions(+), 29 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 8f7f35650949c..615562540b11f 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 @@ -64,6 +64,7 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.ManagedLedgerAlreadyClosedException; import org.apache.bookkeeper.mledger.ManagedLedgerException.ManagedLedgerFencedException; import org.apache.bookkeeper.mledger.ManagedLedgerException.ManagedLedgerTerminatedException; +import org.apache.bookkeeper.mledger.ManagedLedgerException.MetadataNotFoundException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedCursorContainer; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; @@ -955,7 +956,31 @@ public CompletableFuture createSubscription(String subscriptionNam @Override public CompletableFuture unsubscribe(String subscriptionName) { CompletableFuture unsubscribeFuture = new CompletableFuture<>(); + getBrokerService().getManagedLedgerFactory().asyncDelete(TopicName.get(MLPendingAckStore + .getTransactionPendingAckStoreSuffix(topic, subscriptionName)).getPersistenceNamingEncoding(), + new AsyncCallbacks.DeleteLedgerCallback() { + @Override + public void deleteLedgerComplete(Object ctx) { + asyncDeleteCursor(subscriptionName, unsubscribeFuture); + } + + @Override + public void deleteLedgerFailed(ManagedLedgerException exception, Object ctx) { + if (exception instanceof MetadataNotFoundException) { + asyncDeleteCursor(subscriptionName, unsubscribeFuture); + return; + } + + unsubscribeFuture.completeExceptionally(exception); + log.error("[{}][{}] Error deleting subscription pending ack store", + topic, subscriptionName, exception); + } + }, null); + + return unsubscribeFuture; + } + private void asyncDeleteCursor(String subscriptionName, CompletableFuture unsubscribeFuture) { ledger.asyncDeleteCursor(Codec.encode(subscriptionName), new DeleteCursorCallback() { @Override public void deleteCursorComplete(Object ctx) { @@ -970,13 +995,12 @@ public void deleteCursorComplete(Object ctx) { @Override public void deleteCursorFailed(ManagedLedgerException exception, Object ctx) { if (log.isDebugEnabled()) { - log.debug("[{}][{}] Error deleting cursor for subscription", topic, subscriptionName, exception); + log.debug("[{}][{}] Error deleting cursor for subscription", + topic, subscriptionName, exception); } unsubscribeFuture.completeExceptionally(new PersistenceException(exception)); } }, null); - - return unsubscribeFuture; } void removeSubscription(String subscriptionName) { @@ -1083,32 +1107,45 @@ private CompletableFuture delete(boolean failIfHasSubscriptions, unfenceTopicToResume(); deleteFuture.completeExceptionally(ex); } else { - ledger.asyncDelete(new AsyncCallbacks.DeleteLedgerCallback() { - @Override - public void deleteLedgerComplete(Object ctx) { - brokerService.removeTopicFromCache(topic); - - dispatchRateLimiter.ifPresent(DispatchRateLimiter::close); - - subscribeRateLimiter.ifPresent(SubscribeRateLimiter::close); - - brokerService.pulsar().getTopicPoliciesService().clean(TopicName.get(topic)); - log.info("[{}] Topic deleted", topic); - deleteFuture.complete(null); - } - - @Override - public void deleteLedgerFailed(ManagedLedgerException exception, Object ctx) { - if (exception.getCause() instanceof KeeperException.NoNodeException) { - log.info("[{}] Topic is already deleted {}", topic, exception.getMessage()); - deleteLedgerComplete(ctx); - } else { - unfenceTopicToResume(); - log.error("[{}] Error deleting topic", topic, exception); - deleteFuture.completeExceptionally(new PersistenceException(exception)); - } + List> subsDeleteFutures = new ArrayList<>(); + subscriptions.forEach((sub, p) -> subsDeleteFutures.add(unsubscribe(sub))); + + FutureUtil.waitForAll(subsDeleteFutures).whenComplete((f, e) -> { + if (e != null) { + log.error("[{}] Error deleting topic", topic, e); + unfenceTopicToResume(); + deleteFuture.completeExceptionally(e); + } else { + ledger.asyncDelete(new AsyncCallbacks.DeleteLedgerCallback() { + @Override + public void deleteLedgerComplete(Object ctx) { + brokerService.removeTopicFromCache(topic); + + dispatchRateLimiter.ifPresent(DispatchRateLimiter::close); + + subscribeRateLimiter.ifPresent(SubscribeRateLimiter::close); + + brokerService.pulsar().getTopicPoliciesService() + .clean(TopicName.get(topic)); + log.info("[{}] Topic deleted", topic); + deleteFuture.complete(null); + } + + @Override + public void deleteLedgerFailed(ManagedLedgerException exception, Object ctx) { + if (exception.getCause() instanceof KeeperException.NoNodeException) { + log.info("[{}] Topic is already deleted {}", + topic, exception.getMessage()); + deleteLedgerComplete(ctx); + } else { + unfenceTopicToResume(); + log.error("[{}] Error deleting topic", topic, exception); + deleteFuture.completeExceptionally(new PersistenceException(exception)); + } + } + }, null); } - }, null); + }); } }); } else { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/pendingack/PendingAckPersistentTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/pendingack/PendingAckPersistentTest.java index b0d25c08413a5..ae9f12b625175 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/pendingack/PendingAckPersistentTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/pendingack/PendingAckPersistentTest.java @@ -18,6 +18,7 @@ */ package org.apache.pulsar.broker.transaction.pendingack; +import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertTrue; import static org.testng.Assert.fail; import com.google.common.collect.Sets; @@ -36,6 +37,7 @@ import org.apache.pulsar.broker.transaction.TransactionTestBase; import org.apache.pulsar.broker.transaction.pendingack.impl.MLPendingAckStore; import org.apache.pulsar.broker.transaction.pendingack.impl.PendingAckHandleImpl; +import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; @@ -45,6 +47,7 @@ import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.api.transaction.Transaction; import org.apache.pulsar.common.naming.NamespaceName; +import org.apache.pulsar.common.naming.TopicDomain; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.ClusterDataImpl; import org.apache.pulsar.common.policies.data.TenantInfoImpl; @@ -61,6 +64,8 @@ public class PendingAckPersistentTest extends TransactionTestBase { private static final String PENDING_ACK_REPLAY_TOPIC = "persistent://public/txn/pending-ack-replay"; + private static final String NAMESPACE = "public/txn"; + @BeforeMethod public void setup() throws Exception { setBrokerCount(1); @@ -75,7 +80,7 @@ public void setup() throws Exception { admin.topics().createPartitionedTopic(TopicName.TRANSACTION_COORDINATOR_ASSIGN.toString(), 16); admin.tenants().createTenant("public", new TenantInfoImpl(Sets.newHashSet(), Sets.newHashSet(CLUSTER_NAME))); - admin.namespaces().createNamespace("public/txn", 10); + admin.namespaces().createNamespace(NAMESPACE, 10); admin.topics().createNonPartitionedTopic(PENDING_ACK_REPLAY_TOPIC); pulsarClient = PulsarClient.builder() @@ -298,4 +303,67 @@ public void cumulativePendingAckReplayTest() throws Exception { .until(() -> ((PositionImpl) managedCursor.getMarkDeletedPosition()) .compareTo((PositionImpl) managedCursor.getManagedLedger().getLastConfirmedEntry()) == -1); } + + @Test + private void testDeleteSubThenDeletePendingAckManagedLedger() throws Exception { + + String subName = "test-delete"; + + String topic = TopicName.get(TopicDomain.persistent.toString(), + NamespaceName.get(NAMESPACE), "test-delete").toString(); + @Cleanup + Consumer consumer = pulsarClient.newConsumer() + .topic(topic) + .subscriptionName(subName) + .subscriptionType(SubscriptionType.Failover) + .enableBatchIndexAcknowledgment(true) + .subscribe(); + + consumer.close(); + + admin.topics().deleteSubscription(topic, subName); + + List topics = admin.namespaces().getTopics(NAMESPACE); + + assertFalse(topics.contains(MLPendingAckStore.getTransactionPendingAckStoreSuffix(topic, subName))); + + assertTrue(topics.contains(topic)); + } + + @Test + private void testDeleteTopicThenDeletePendingAckManagedLedger() throws Exception { + + String subName1 = "test-delete"; + String subName2 = "test-delete"; + + String topic = TopicName.get(TopicDomain.persistent.toString(), + NamespaceName.get(NAMESPACE), "test-delete").toString(); + @Cleanup + Consumer consumer1 = pulsarClient.newConsumer() + .topic(topic) + .subscriptionName(subName1) + .subscriptionType(SubscriptionType.Failover) + .enableBatchIndexAcknowledgment(true) + .subscribe(); + + consumer1.close(); + + @Cleanup + Consumer consumer2 = pulsarClient.newConsumer() + .topic(topic) + .subscriptionName(subName2) + .subscriptionType(SubscriptionType.Failover) + .enableBatchIndexAcknowledgment(true) + .subscribe(); + + consumer2.close(); + + admin.topics().delete(topic); + + List topics = admin.namespaces().getTopics(NAMESPACE); + + assertFalse(topics.contains(MLPendingAckStore.getTransactionPendingAckStoreSuffix(topic, subName1))); + assertFalse(topics.contains(MLPendingAckStore.getTransactionPendingAckStoreSuffix(topic, subName2))); + assertFalse(topics.contains(topic)); + } } From 86f39f290aa714d83b5112f7cacd122e0affac4c Mon Sep 17 00:00:00 2001 From: congbo Date: Thu, 1 Jul 2021 14:08:35 +0800 Subject: [PATCH 2/6] fix code style --- .../apache/pulsar/broker/service/persistent/PersistentTopic.java | 1 + 1 file changed, 1 insertion(+) 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 33878a072172c..9d0bb2f428715 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 @@ -108,6 +108,7 @@ import org.apache.pulsar.broker.stats.ReplicationMetrics; import org.apache.pulsar.broker.transaction.buffer.TransactionBuffer; import org.apache.pulsar.broker.transaction.buffer.impl.TransactionBufferDisable; +import org.apache.pulsar.broker.transaction.pendingack.impl.MLPendingAckStore; import org.apache.pulsar.client.admin.LongRunningProcessStatus; import org.apache.pulsar.client.admin.OffloadProcessStatus; import org.apache.pulsar.client.api.MessageId; From 5c6b08644cdb799825341ffb9b55ed9869fbee1a Mon Sep 17 00:00:00 2001 From: congbo Date: Fri, 9 Jul 2021 00:05:40 +0800 Subject: [PATCH 3/6] Fix code style --- .../pulsar/broker/service/persistent/PersistentTopic.java | 4 ++-- 1 file changed, 2 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 a3d41fd5d66e1..d259f8617e056 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 @@ -1120,8 +1120,8 @@ public void deleteLedgerComplete(Object ctx) { @Override public void deleteLedgerFailed(ManagedLedgerException exception, Object ctx) { - if (exception.getCause() instanceof - MetadataStoreException.NotFoundException) { + if (exception.getCause() + instanceof MetadataStoreException.NotFoundException) { log.info("[{}] Topic is already deleted {}", topic, exception.getMessage()); deleteLedgerComplete(ctx); From 17653506d337c72c7cc85f9a0cd1a500cca5d870 Mon Sep 17 00:00:00 2001 From: congbo Date: Fri, 9 Jul 2021 11:17:05 +0800 Subject: [PATCH 4/6] fix some test --- .../transaction/pendingack/PendingAckPersistentTest.java | 3 +++ .../tests/integration/functions/PulsarFunctionsTest.java | 8 ++++++-- 2 files changed, 9 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/pendingack/PendingAckPersistentTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/pendingack/PendingAckPersistentTest.java index ae9f12b625175..506e9cfc11d8d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/pendingack/PendingAckPersistentTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/pendingack/PendingAckPersistentTest.java @@ -51,6 +51,7 @@ import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.ClusterDataImpl; import org.apache.pulsar.common.policies.data.TenantInfoImpl; +import org.apache.pulsar.common.policies.data.TopicStats; import org.awaitility.Awaitility; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -325,6 +326,8 @@ private void testDeleteSubThenDeletePendingAckManagedLedger() throws Exception { List topics = admin.namespaces().getTopics(NAMESPACE); + TopicStats topicStats = admin.topics().getStats(topic, false); + assertFalse(topics.contains(MLPendingAckStore.getTransactionPendingAckStoreSuffix(topic, subName))); assertTrue(topics.contains(topic)); diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java index 87357a04cf6fb..3a501d1868f02 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java @@ -729,8 +729,12 @@ protected void testExclamationFunction(Runtime runtime, // get function info getFunctionInfoNotFound(functionName); - // make sure subscriptions are cleanup - checkSubscriptionsCleanup(inputTopicName); + final String topic = inputTopicName; + Awaitility.await().until(() -> { + // make sure subscriptions are cleanup + checkSubscriptionsCleanup(topic); + return true; + }); } From 8e2770e9f71454f9bd641d75bbe236018a2fc309 Mon Sep 17 00:00:00 2001 From: congbo Date: Fri, 6 Aug 2021 16:24:02 +0800 Subject: [PATCH 5/6] Fix some test --- .../java/org/apache/pulsar/broker/service/ServerCnxTest.java | 4 ++-- .../transaction/buffer/TransactionBufferClientTest.java | 5 ++++- 2 files changed, 6 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java index f16f7bbf85bbf..d670f504378e6 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java @@ -772,7 +772,7 @@ public void testDuplicateConcurrentSubscribeCommand() throws Exception { // Create producer second time clientCommand = Commands.newSubscribe(successTopicName, // - successSubName, 1 /* consumer id */, 1 /* request id */, SubType.Exclusive, 0, + successSubName, 2 /* consumer id */, 1 /* request id */, SubType.Exclusive, 0, "test" /* consumer name */, 0 /* avoid reseting cursor */); channel.writeInbound(clientCommand); @@ -780,7 +780,7 @@ public void testDuplicateConcurrentSubscribeCommand() throws Exception { Object response = getResponse(); assertTrue(response instanceof CommandError, "Response is not CommandError but " + response); CommandError error = (CommandError) response; - assertEquals(error.getError(), ServerError.ServiceNotReady); + assertEquals(error.getError(), ServerError.ConsumerBusy); }); channel.finish(); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/buffer/TransactionBufferClientTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/buffer/TransactionBufferClientTest.java index 2ad16a79cd583..ba73034cb9cf9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/buffer/TransactionBufferClientTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/buffer/TransactionBufferClientTest.java @@ -46,6 +46,7 @@ import org.apache.pulsar.broker.transaction.buffer.impl.TransactionBufferClientImpl; import org.apache.pulsar.broker.transaction.buffer.impl.TransactionBufferHandlerImpl; import org.apache.pulsar.broker.transaction.coordinator.TransactionMetaStoreTestBase; +import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.transaction.TransactionBufferClient; import org.apache.pulsar.client.api.transaction.TransactionBufferClientException; @@ -90,9 +91,11 @@ protected void afterSetup() throws Exception { pulsarAdmins[0].tenants().createTenant("public", new TenantInfoImpl(Sets.newHashSet(), Sets.newHashSet("my-cluster"))); pulsarAdmins[0].namespaces().createNamespace(namespace, 10); pulsarAdmins[0].topics().createPartitionedTopic(partitionedTopicName.getPartitionedTopicName(), partitions); + String subName = "test"; + pulsarAdmins[0].topics().createSubscription(partitionedTopicName.getPartitionedTopicName(), subName, MessageId.latest); pulsarClient.newConsumer() .topic(partitionedTopicName.getPartitionedTopicName()) - .subscriptionName("test").subscribe(); + .subscriptionName(subName).subscribe(); tbClient = TransactionBufferClientImpl.create( ((PulsarClientImpl) pulsarClient), new HashedWheelTimer(new DefaultThreadFactory("transaction-buffer"))); From 7bc4530eb807b3930f4243e518db010b329aa799 Mon Sep 17 00:00:00 2001 From: congbo Date: Sat, 7 Aug 2021 11:41:34 +0800 Subject: [PATCH 6/6] Fix some mock test --- .../apache/pulsar/broker/service/PersistentTopicTest.java | 6 ++++++ 1 file changed, 6 insertions(+) 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 9b74331450c9c..58d3b24cb1a53 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 @@ -182,6 +182,12 @@ public void setup() throws Exception { mlFactoryMock = mock(ManagedLedgerFactory.class); doReturn(mlFactoryMock).when(pulsar).getManagedLedgerFactory(); + doAnswer(invocation -> { + DeleteLedgerCallback deleteLedgerCallback = invocation.getArgument(1); + deleteLedgerCallback.deleteLedgerComplete(null); + return null; + }).when(mlFactoryMock).asyncDelete(any(), any(), any()); + ZooKeeper mockZk = createMockZooKeeper(); doReturn(mockZk).when(pulsar).getZkClient(); doReturn(createMockBookKeeper(executor))