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 681cdb0fdf7b5..b82e73fee91f9 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 @@ -3636,6 +3636,21 @@ private CompletableFuture updatePartitionedTopic(TopicName topicName, int p -> new PartitionedTopicMetadata(numPartitions)); updatePartition.complete(null); } catch (Exception e) { + getPartitionedTopicMetadataAsync(topicName, false, false).thenAccept(metadata -> { + int oldPartition = metadata.partitions; + for (int i = oldPartition; i < numPartitions; i++) { + String managedLedgerPath = ZkAdminPaths.managedLedgerPath(topicName.getPartition(i)); + namespaceResources().getPartitionedTopicResources() + .deleteAsync(managedLedgerPath).exceptionally(ex1 -> { + log.warn("[{}] Failed to clean up managedLedger znode {}", clientAppId(), + managedLedgerPath, ex1.getCause()); + return null; + }); + } + }).exceptionally(ex -> { + log.warn("[{}] Failed to clean up managedLedger znode", clientAppId(), ex.getCause()); + return null; + }); updatePartition.completeExceptionally(e); } }).exceptionally(ex -> {