diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index e455b9ce81c63..511895354abef 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -1150,9 +1150,10 @@ void handleFailure(Throwable throwable) { } @Override - public DeleteTopicsResult deleteTopics(final Collection topicNames, + public DeleteTopicsResult deleteTopics(Collection topicNames, DeleteTopicsOptions options) { final Map> topicFutures = new HashMap<>(topicNames.size()); + final List validTopicNames = new ArrayList<>(topicNames.size()); for (String topicName : topicNames) { if (topicNameIsUnrepresentable(topicName)) { KafkaFutureImpl future = new KafkaFutureImpl<>(); @@ -1161,6 +1162,7 @@ public DeleteTopicsResult deleteTopics(final Collection topicNames, topicFutures.put(topicName, future); } else if (!topicFutures.containsKey(topicName)) { topicFutures.put(topicName, new KafkaFutureImpl()); + validTopicNames.add(topicName); } } final long now = time.milliseconds(); @@ -1169,7 +1171,7 @@ public DeleteTopicsResult deleteTopics(final Collection topicNames, @Override AbstractRequest.Builder createRequest(int timeoutMs) { - return new DeleteTopicsRequest.Builder(new HashSet<>(topicNames), timeoutMs); + return new DeleteTopicsRequest.Builder(new HashSet<>(validTopicNames), timeoutMs); } @Override @@ -1204,7 +1206,7 @@ void handleFailure(Throwable throwable) { completeAllExceptionally(topicFutures.values(), throwable); } }; - if (!topicNames.isEmpty()) { + if (!validTopicNames.isEmpty()) { runnable.call(call, now); } return new DeleteTopicsResult(new HashMap>(topicFutures)); diff --git a/clients/src/test/java/org/apache/kafka/clients/MockClient.java b/clients/src/test/java/org/apache/kafka/clients/MockClient.java index 60af9bcc84937..a73175c995430 100644 --- a/clients/src/test/java/org/apache/kafka/clients/MockClient.java +++ b/clients/src/test/java/org/apache/kafka/clients/MockClient.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.clients; +import java.util.concurrent.atomic.AtomicInteger; import org.apache.kafka.common.Cluster; import org.apache.kafka.common.Node; import org.apache.kafka.common.errors.AuthenticationException; @@ -88,7 +89,7 @@ public FutureResponse(Node node, private final Queue metadataUpdates = new ArrayDeque<>(); private volatile NodeApiVersions nodeApiVersions = NodeApiVersions.create(); private volatile int numBlockingWakeups = 0; - + private final AtomicInteger totalRequestCount = new AtomicInteger(0); public MockClient(Time time) { this(time, null); } @@ -394,6 +395,7 @@ public void reset() { futureResponses.clear(); metadataUpdates.clear(); authenticationErrors.clear(); + totalRequestCount.set(0); } public boolean hasPendingMetadataUpdates() { @@ -461,6 +463,7 @@ public ClientRequest newClientRequest(String nodeId, AbstractRequest.Builder @Override public ClientRequest newClientRequest(String nodeId, AbstractRequest.Builder requestBuilder, long createdTimeMs, boolean expectResponse, RequestCompletionHandler callback) { + totalRequestCount.incrementAndGet(); return new ClientRequest(nodeId, requestBuilder, 0, "mockClientId", createdTimeMs, expectResponse, callback); } @@ -503,4 +506,9 @@ private static class MetadataUpdate { this.expectMatchRefreshTopics = expectMatchRefreshTopics; } } + + // visible for testing + public int totalRequestCount() { + return totalRequestCount.get(); + } } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index f08a99b6ddc4b..0d4dee65c1c6c 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -230,22 +230,26 @@ public void testInvalidTopicNames() throws Exception { for (String sillyTopicName : sillyTopicNames) { assertFutureError(deleteFutures.get(sillyTopicName), InvalidTopicException.class); } + assertEquals(0, env.kafkaClient().totalRequestCount()); Map> describeFutures = env.adminClient().describeTopics(sillyTopicNames).values(); for (String sillyTopicName : sillyTopicNames) { assertFutureError(describeFutures.get(sillyTopicName), InvalidTopicException.class); } + assertEquals(0, env.kafkaClient().totalRequestCount()); List newTopics = new ArrayList<>(); for (String sillyTopicName : sillyTopicNames) { newTopics.add(new NewTopic(sillyTopicName, 1, (short) 1)); } + Map> createFutures = env.adminClient().createTopics(newTopics).values(); for (String sillyTopicName : sillyTopicNames) { assertFutureError(createFutures .get(sillyTopicName), InvalidTopicException.class); } + assertEquals(0, env.kafkaClient().totalRequestCount()); } }