From 3d554d1c5beae04911609d733d0b69005a6de49e Mon Sep 17 00:00:00 2001 From: fanjianye Date: Tue, 19 Apr 2022 17:43:52 +0800 Subject: [PATCH 1/2] After catch the "Subscription already exists" error, do updatePartitionedTopicAsync() --- .../admin/impl/PersistentTopicsBase.java | 21 ++++++++++++++++++- .../pulsar/broker/admin/AdminApi2Test.java | 8 ++++++- 2 files changed, 27 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 3c3f622266c56..91583546edad1 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -4052,7 +4052,26 @@ private CompletableFuture updatePartitionedTopic(TopicName topicName, int return future; }).thenAccept(__ -> result.complete(null)).exceptionally(ex -> { if (force && ex.getCause() instanceof PulsarAdminException.ConflictException) { - result.complete(null); + CompletableFuture future = namespaceResources().getPartitionedTopicResources() + .updatePartitionedTopicAsync(topicName, p -> new PartitionedTopicMetadata(numPartitions)); + future.thenAccept(__ -> result.complete(null)).exceptionally(ex2 -> { + // If the update operation fails, clean up the partitions that were created + getPartitionedTopicMetadataAsync(topicName, false, false).thenAccept(metadata -> { + int oldPartition = metadata.partitions; + for (int i = oldPartition; i < numPartitions; i++) { + topicResources().deletePersistentTopicAsync(topicName.getPartition(i)).exceptionally(ex1 -> { + log.warn("[{}] Failed to clean up managedLedger {}", clientAppId(), topicName, + ex1.getCause()); + return null; + }); + } + }).exceptionally(e -> { + log.warn("[{}] Failed to clean up managedLedger", topicName, e); + return null; + }); + result.completeExceptionally(ex2); + return null; + }); return null; } result.completeExceptionally(ex); 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 a6b89112634de..b683768efa08d 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 @@ -2350,7 +2350,13 @@ public void testFailedUpdatePartitionedTopic() throws Exception { } admin.topics().updatePartitionedTopic(partitionedTopicName, newPartitions, false, true); // validate subscription is created for new partition. - assertNotNull(admin.topics().getStats(partitionedTopicName + "-partition-" + 6).getSubscriptions().get(subName1)); + for (int i = startPartitions; i < newPartitions; i++) { + assertNotNull( + admin.topics().getStats(partitionedTopicName + "-partition-" + i).getSubscriptions().get(subName1)); + } + + // validate update partition is success + assertEquals(admin.topics().getPartitionedTopicMetadata(partitionedTopicName).partitions, newPartitions); } @Test(dataProvider = "topicType") From e90def9be7c9811705190492f83f0761e138ad77 Mon Sep 17 00:00:00 2001 From: fanjianye Date: Thu, 21 Apr 2022 15:32:42 +0800 Subject: [PATCH 2/2] no need to clean up znode when updatePartitionedMetadata failed --- .../admin/impl/PersistentTopicsBase.java | 43 +++---------------- 1 file changed, 5 insertions(+), 38 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 91583546edad1..7993600c3a49b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -4029,47 +4029,14 @@ private PersistentReplicator getReplicatorReference(String replName, PersistentT private CompletableFuture updatePartitionedTopic(TopicName topicName, int numPartitions, boolean force) { CompletableFuture result = new CompletableFuture<>(); - createSubscriptions(topicName, numPartitions).thenCompose(__ -> { - CompletableFuture future = namespaceResources().getPartitionedTopicResources() - .updatePartitionedTopicAsync(topicName, p -> new PartitionedTopicMetadata(numPartitions)); - future.exceptionally(ex -> { - // If the update operation fails, clean up the partitions that were created - getPartitionedTopicMetadataAsync(topicName, false, false).thenAccept(metadata -> { - int oldPartition = metadata.partitions; - for (int i = oldPartition; i < numPartitions; i++) { - topicResources().deletePersistentTopicAsync(topicName.getPartition(i)).exceptionally(ex1 -> { - log.warn("[{}] Failed to clean up managedLedger {}", clientAppId(), topicName, - ex1.getCause()); - return null; - }); - } - }).exceptionally(e -> { - log.warn("[{}] Failed to clean up managedLedger", topicName, e); - return null; - }); - return null; - }); - return future; - }).thenAccept(__ -> result.complete(null)).exceptionally(ex -> { + createSubscriptions(topicName, numPartitions).thenCompose(__ -> namespaceResources().getPartitionedTopicResources() + .updatePartitionedTopicAsync(topicName, p -> new PartitionedTopicMetadata(numPartitions)) + ).thenAccept(__ -> result.complete(null)).exceptionally(ex -> { if (force && ex.getCause() instanceof PulsarAdminException.ConflictException) { CompletableFuture future = namespaceResources().getPartitionedTopicResources() .updatePartitionedTopicAsync(topicName, p -> new PartitionedTopicMetadata(numPartitions)); - future.thenAccept(__ -> result.complete(null)).exceptionally(ex2 -> { - // If the update operation fails, clean up the partitions that were created - getPartitionedTopicMetadataAsync(topicName, false, false).thenAccept(metadata -> { - int oldPartition = metadata.partitions; - for (int i = oldPartition; i < numPartitions; i++) { - topicResources().deletePersistentTopicAsync(topicName.getPartition(i)).exceptionally(ex1 -> { - log.warn("[{}] Failed to clean up managedLedger {}", clientAppId(), topicName, - ex1.getCause()); - return null; - }); - } - }).exceptionally(e -> { - log.warn("[{}] Failed to clean up managedLedger", topicName, e); - return null; - }); - result.completeExceptionally(ex2); + future.thenAccept(__ -> result.complete(null)).exceptionally(ex1 -> { + result.completeExceptionally(ex1); return null; }); return null;