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..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,31 +4029,17 @@ 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); + 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(ex1 -> { + result.completeExceptionally(ex1); return null; }); return null; - }); - return future; - }).thenAccept(__ -> result.complete(null)).exceptionally(ex -> { - if (force && ex.getCause() instanceof PulsarAdminException.ConflictException) { - result.complete(null); - return null; } result.completeExceptionally(ex); return null; 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")