From e699285022f3647be29bb976d4a31544b79aebbc Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Thu, 16 Dec 2021 12:28:22 +0100 Subject: [PATCH 1/3] AdminManager: handle properly async deleteParitionedTopic --- .../pulsar/handlers/kop/AdminManager.java | 30 ++++++++++++------- 1 file changed, 19 insertions(+), 11 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java index 333271394d..32b8eabfae 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java @@ -51,6 +51,7 @@ import org.apache.kafka.common.requests.DescribeConfigsResponse; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminException; +import org.apache.pulsar.common.util.FutureUtil; @Slf4j class AdminManager { @@ -63,6 +64,7 @@ class AdminManager { private final PulsarAdmin admin; private final int defaultNumPartitions; + private volatile Map> brokersCache = Maps.newHashMap(); private final ReentrantReadWriteLock brokersCacheLock = new ReentrantReadWriteLock(); @@ -212,7 +214,7 @@ CompletableFuture> describeC future.complete(new DescribeConfigsResponse.Config(ApiError.NONE, dummyConfig)); break; default: - throw new InvalidRequestException("Unsupported resource type: " + resource.type()); + return FutureUtil.failedFuture(new InvalidRequestException("Unsupported resource type: " + resource.type())); } return future; } catch (Exception e) { @@ -222,6 +224,10 @@ CompletableFuture> describeC })); CompletableFuture> resultFuture = new CompletableFuture<>(); CompletableFuture.allOf(futureMap.values().toArray(new CompletableFuture[0])).whenComplete((ignored, e) -> { + if (e != null) { + resultFuture.completeExceptionally(e); + return; + } resultFuture.complete(futureMap.entrySet().stream().collect( Collectors.toMap(Map.Entry::getKey, entry -> entry.getValue().getNow(null)) )); @@ -240,14 +246,17 @@ private DescribeConfigsResponse.ConfigEntry buildDummyEntryConfig(String configN public void deleteTopic(String topicToDelete, Consumer successConsumer, Consumer errorConsumer) { - try { - admin.topics().deletePartitionedTopic(topicToDelete); - successConsumer.accept(topicToDelete); - log.info("delete topic {} successfully.", topicToDelete); - } catch (PulsarAdminException e) { - log.error("delete topic {} failed, exception: ", topicToDelete, e); - errorConsumer.accept(topicToDelete); - } + admin.topics() + .deletePartitionedTopicAsync(topicToDelete) + .thenRun(() -> { + log.info("delete topic {} successfully.", topicToDelete); + successConsumer.accept(topicToDelete); + }) + .exceptionally((e -> { + log.error("delete topic {} failed, exception: ", topicToDelete, e); + errorConsumer.accept(topicToDelete); + return null; + })); } public void truncateTopic(String topicToDelete, @@ -278,8 +287,6 @@ public void truncateTopic(String topicToDelete, } - - CompletableFuture> createPartitionsAsync(Map createInfo, int timeoutMs, String namespacePrefix) { @@ -434,4 +441,5 @@ public void setBrokers(Map> newBrokers) { brokersCacheLock.writeLock().unlock(); } } + } From 4557c8211a8dcc41ee05c9dbf740d3646645ff75 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Thu, 16 Dec 2021 13:05:33 +0100 Subject: [PATCH 2/3] checkstyle --- .../java/io/streamnative/pulsar/handlers/kop/AdminManager.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java index 32b8eabfae..86f2a7c252 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java @@ -214,7 +214,8 @@ CompletableFuture> describeC future.complete(new DescribeConfigsResponse.Config(ApiError.NONE, dummyConfig)); break; default: - return FutureUtil.failedFuture(new InvalidRequestException("Unsupported resource type: " + resource.type())); + return FutureUtil.failedFuture( + new InvalidRequestException("Unsupported resource type: " + resource.type())); } return future; } catch (Exception e) { From a6db7852984ad967b606c27263704493cab146a9 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Tue, 21 Dec 2021 10:41:21 +0100 Subject: [PATCH 3/3] address comments --- .../io/streamnative/pulsar/handlers/kop/AdminManager.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java index 86f2a7c252..3ec3fb964e 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java @@ -51,7 +51,6 @@ import org.apache.kafka.common.requests.DescribeConfigsResponse; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminException; -import org.apache.pulsar.common.util.FutureUtil; @Slf4j class AdminManager { @@ -214,8 +213,11 @@ CompletableFuture> describeC future.complete(new DescribeConfigsResponse.Config(ApiError.NONE, dummyConfig)); break; default: - return FutureUtil.failedFuture( - new InvalidRequestException("Unsupported resource type: " + resource.type())); + return CompletableFuture.completedFuture(new DescribeConfigsResponse.Config( + ApiError.fromThrowable( + new InvalidRequestException("Unsupported resource type: " + + resource.type())), + Collections.emptyList())); } return future; } catch (Exception e) {