From 04e2a1e2dca3f0963b1ffed368a9669ae3ceccbb Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Tue, 17 Aug 2021 17:39:46 +0800 Subject: [PATCH 1/6] clean zk managed-ledgers after failed update partition --- .../pulsar/broker/admin/impl/PersistentTopicsBase.java | 9 +++++++++ 1 file changed, 9 insertions(+) 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..e7d90a2bd63d1 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 @@ -3628,6 +3628,8 @@ private PersistentReplicator getReplicatorReference(String replName, PersistentT private CompletableFuture updatePartitionedTopic(TopicName topicName, int numPartitions) { final String path = ZkAdminPaths.partitionedTopicPath(topicName); + PartitionedTopicMetadata metadata = getPartitionedTopicMetadata(topicName, false, false); + int oldPartition = metadata.partitions; CompletableFuture updatePartition = new CompletableFuture<>(); createSubscriptions(topicName, numPartitions).thenAccept(res -> { @@ -3636,6 +3638,13 @@ private CompletableFuture updatePartitionedTopic(TopicName topicName, int p -> new PartitionedTopicMetadata(numPartitions)); updatePartition.complete(null); } catch (Exception e) { + for(int i=oldPartition; i { From bc771fe8c0f606df1609e6c7eabdea7e99630afa Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Wed, 18 Aug 2021 13:35:11 +0800 Subject: [PATCH 2/6] support async method --- .../broker/admin/impl/PersistentTopicsBase.java | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 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 e7d90a2bd63d1..50605c272d078 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 @@ -3628,8 +3628,6 @@ private PersistentReplicator getReplicatorReference(String replName, PersistentT private CompletableFuture updatePartitionedTopic(TopicName topicName, int numPartitions) { final String path = ZkAdminPaths.partitionedTopicPath(topicName); - PartitionedTopicMetadata metadata = getPartitionedTopicMetadata(topicName, false, false); - int oldPartition = metadata.partitions; CompletableFuture updatePartition = new CompletableFuture<>(); createSubscriptions(topicName, numPartitions).thenAccept(res -> { @@ -3638,13 +3636,15 @@ private CompletableFuture updatePartitionedTopic(TopicName topicName, int p -> new PartitionedTopicMetadata(numPartitions)); updatePartition.complete(null); } catch (Exception e) { - for(int i=oldPartition; i { + int oldPartition = metadata.partitions; + for(int i = oldPartition; i < numPartitions; i++){ + namespaceResources().getPartitionedTopicResources().deleteAsync(ZkAdminPaths.managedLedgerPath(topicName.getPartition(i))); } - } + }).exceptionally(ex -> { + updatePartition.completeExceptionally(e); + return null; + }); updatePartition.completeExceptionally(e); } }).exceptionally(ex -> { From 8dc16b78eb9be72e4c5c87abba5977e1e768ca33 Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Wed, 18 Aug 2021 15:42:05 +0800 Subject: [PATCH 3/6] fix code style --- .../pulsar/broker/admin/impl/PersistentTopicsBase.java | 9 +++++++-- 1 file changed, 7 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 50605c272d078..a41fdb46bf47a 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 @@ -3638,8 +3638,13 @@ private CompletableFuture updatePartitionedTopic(TopicName topicName, int } catch (Exception e) { getPartitionedTopicMetadataAsync(topicName, false, false).thenAccept(metadata -> { int oldPartition = metadata.partitions; - for(int i = oldPartition; i < numPartitions; i++){ - namespaceResources().getPartitionedTopicResources().deleteAsync(ZkAdminPaths.managedLedgerPath(topicName.getPartition(i))); + for (int i = oldPartition; i < numPartitions; i++) { + String managedLedgerPath = ZkAdminPaths.managedLedgerPath(topicName.getPartition(i)); + try { + namespaceResources().getPartitionedTopicResources().deleteAsync(managedLedgerPath); + } catch (Exception cleanZnodeException) { + log.error("Failed to clean managedLedger znode {}", managedLedgerPath, cleanZnodeException); + } } }).exceptionally(ex -> { updatePartition.completeExceptionally(e); From 38c9ddd6c82d5cfe2c3b6148631a808d70b51aec Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Wed, 18 Aug 2021 15:59:44 +0800 Subject: [PATCH 4/6] add error log delete when failed to znode --- .../pulsar/broker/admin/impl/PersistentTopicsBase.java | 9 ++++----- 1 file changed, 4 insertions(+), 5 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 a41fdb46bf47a..0eac927f8b0e9 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 @@ -3640,11 +3640,10 @@ private CompletableFuture updatePartitionedTopic(TopicName topicName, int int oldPartition = metadata.partitions; for (int i = oldPartition; i < numPartitions; i++) { String managedLedgerPath = ZkAdminPaths.managedLedgerPath(topicName.getPartition(i)); - try { - namespaceResources().getPartitionedTopicResources().deleteAsync(managedLedgerPath); - } catch (Exception cleanZnodeException) { - log.error("Failed to clean managedLedger znode {}", managedLedgerPath, cleanZnodeException); - } + namespaceResources().getPartitionedTopicResources().deleteAsync(managedLedgerPath).exceptionally(ex1 -> { + log.error("[{}] Failed to delete managedLedger znode {}", clientAppId(), managedLedgerPath, ex1.getCause()); + return null; + }); } }).exceptionally(ex -> { updatePartition.completeExceptionally(e); From ddee12f081ebe90a7ace0ab9e66ad72e1f43668f Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Wed, 18 Aug 2021 16:13:52 +0800 Subject: [PATCH 5/6] fix code style --- .../pulsar/broker/admin/impl/PersistentTopicsBase.java | 6 ++++-- 1 file changed, 4 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 0eac927f8b0e9..2e4033c655cb6 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 @@ -3640,8 +3640,10 @@ private CompletableFuture updatePartitionedTopic(TopicName topicName, int 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.error("[{}] Failed to delete managedLedger znode {}", clientAppId(), managedLedgerPath, ex1.getCause()); + namespaceResources().getPartitionedTopicResources() + .deleteAsync(managedLedgerPath).exceptionally(ex1 -> { + log.error("[{}] Failed to delete managedLedger znode {}", clientAppId(), + managedLedgerPath, ex1.getCause()); return null; }); } From 16f55e821711355a0a991ac4c914859b00dea6fc Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Thu, 19 Aug 2021 13:51:07 +0800 Subject: [PATCH 6/6] fix log level to warn and delete repeated complete future --- .../apache/pulsar/broker/admin/impl/PersistentTopicsBase.java | 4 ++-- 1 file changed, 2 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 2e4033c655cb6..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 @@ -3642,13 +3642,13 @@ private CompletableFuture updatePartitionedTopic(TopicName topicName, int String managedLedgerPath = ZkAdminPaths.managedLedgerPath(topicName.getPartition(i)); namespaceResources().getPartitionedTopicResources() .deleteAsync(managedLedgerPath).exceptionally(ex1 -> { - log.error("[{}] Failed to delete managedLedger znode {}", clientAppId(), + log.warn("[{}] Failed to clean up managedLedger znode {}", clientAppId(), managedLedgerPath, ex1.getCause()); return null; }); } }).exceptionally(ex -> { - updatePartition.completeExceptionally(e); + log.warn("[{}] Failed to clean up managedLedger znode", clientAppId(), ex.getCause()); return null; }); updatePartition.completeExceptionally(e);