From 7199d9cf47cee7716c9314aa2a745b2a0d6b69bd Mon Sep 17 00:00:00 2001 From: Chia-Ping Tsai Date: Fri, 23 Mar 2018 17:33:54 +0800 Subject: [PATCH 1/4] =?UTF-8?q?Minor:=20Don=E2=80=99t=20send=20the=20Delet?= =?UTF-8?q?eTopicsRequest=20for=20the=20invalid=20topic=20names?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../org/apache/kafka/clients/admin/KafkaAdminClient.java | 8 +++++--- .../test/java/org/apache/kafka/clients/MockClient.java | 9 ++++++++- .../apache/kafka/clients/admin/KafkaAdminClientTest.java | 7 +++++++ 3 files changed, 20 insertions(+), 4 deletions(-) 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..5ea91eea0bd08 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 requestCount = new AtomicInteger(0); public MockClient(Time time) { this(time, null); } @@ -461,6 +462,7 @@ public ClientRequest newClientRequest(String nodeId, AbstractRequest.Builder @Override public ClientRequest newClientRequest(String nodeId, AbstractRequest.Builder requestBuilder, long createdTimeMs, boolean expectResponse, RequestCompletionHandler callback) { + requestCount.incrementAndGet(); return new ClientRequest(nodeId, requestBuilder, 0, "mockClientId", createdTimeMs, expectResponse, callback); } @@ -503,4 +505,9 @@ private static class MetadataUpdate { this.expectMatchRefreshTopics = expectMatchRefreshTopics; } } + + // visible for testing + public int requestCount() { + return requestCount.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..04ff40b46930c 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 @@ -225,27 +225,34 @@ public void testInvalidTopicNames() throws Exception { env.kafkaClient().setNode(env.cluster().controller()); List sillyTopicNames = Arrays.asList(new String[] {"", null}); + int originRequestCount = env.kafkaClient().requestCount(); Map> deleteFutures = env.adminClient().deleteTopics(sillyTopicNames).values(); for (String sillyTopicName : sillyTopicNames) { assertFutureError(deleteFutures.get(sillyTopicName), InvalidTopicException.class); } + assertEquals(originRequestCount, env.kafkaClient().requestCount()); + originRequestCount = env.kafkaClient().requestCount(); Map> describeFutures = env.adminClient().describeTopics(sillyTopicNames).values(); for (String sillyTopicName : sillyTopicNames) { assertFutureError(describeFutures.get(sillyTopicName), InvalidTopicException.class); } + assertEquals(originRequestCount, env.kafkaClient().requestCount()); List newTopics = new ArrayList<>(); for (String sillyTopicName : sillyTopicNames) { newTopics.add(new NewTopic(sillyTopicName, 1, (short) 1)); } + + originRequestCount = env.kafkaClient().requestCount(); Map> createFutures = env.adminClient().createTopics(newTopics).values(); for (String sillyTopicName : sillyTopicNames) { assertFutureError(createFutures .get(sillyTopicName), InvalidTopicException.class); } + assertEquals(originRequestCount, env.kafkaClient().requestCount()); } } From f8af8900744502b5b26c3f0694fd2631f01fa09e Mon Sep 17 00:00:00 2001 From: Chia-Ping Tsai Date: Wed, 28 Mar 2018 15:01:35 +0800 Subject: [PATCH 2/4] address Gustafson's comment - the request count is expected to be 0 --- .../apache/kafka/clients/admin/KafkaAdminClientTest.java | 9 +++------ 1 file changed, 3 insertions(+), 6 deletions(-) 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 04ff40b46930c..f382ea96ec163 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 @@ -225,34 +225,31 @@ public void testInvalidTopicNames() throws Exception { env.kafkaClient().setNode(env.cluster().controller()); List sillyTopicNames = Arrays.asList(new String[] {"", null}); - int originRequestCount = env.kafkaClient().requestCount(); Map> deleteFutures = env.adminClient().deleteTopics(sillyTopicNames).values(); for (String sillyTopicName : sillyTopicNames) { assertFutureError(deleteFutures.get(sillyTopicName), InvalidTopicException.class); } - assertEquals(originRequestCount, env.kafkaClient().requestCount()); + assertEquals(0, env.kafkaClient().requestCount()); - originRequestCount = env.kafkaClient().requestCount(); Map> describeFutures = env.adminClient().describeTopics(sillyTopicNames).values(); for (String sillyTopicName : sillyTopicNames) { assertFutureError(describeFutures.get(sillyTopicName), InvalidTopicException.class); } - assertEquals(originRequestCount, env.kafkaClient().requestCount()); + assertEquals(0, env.kafkaClient().requestCount()); List newTopics = new ArrayList<>(); for (String sillyTopicName : sillyTopicNames) { newTopics.add(new NewTopic(sillyTopicName, 1, (short) 1)); } - originRequestCount = env.kafkaClient().requestCount(); Map> createFutures = env.adminClient().createTopics(newTopics).values(); for (String sillyTopicName : sillyTopicNames) { assertFutureError(createFutures .get(sillyTopicName), InvalidTopicException.class); } - assertEquals(originRequestCount, env.kafkaClient().requestCount()); + assertEquals(0, env.kafkaClient().requestCount()); } } From 81130c22fd99d699eddcf8cfeced4a12776c0949 Mon Sep 17 00:00:00 2001 From: Chia-Ping Tsai Date: Thu, 29 Mar 2018 01:33:05 +0800 Subject: [PATCH 3/4] address Gustafson's comment - rename requestCount to totalRequestCount and reset totalRequestCount in MockClient.reset --- .../src/test/java/org/apache/kafka/clients/MockClient.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) 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 5ea91eea0bd08..f3123dc43181d 100644 --- a/clients/src/test/java/org/apache/kafka/clients/MockClient.java +++ b/clients/src/test/java/org/apache/kafka/clients/MockClient.java @@ -89,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 requestCount = new AtomicInteger(0); + private final AtomicInteger totalRequestCount = new AtomicInteger(0); public MockClient(Time time) { this(time, null); } @@ -395,6 +395,7 @@ public void reset() { futureResponses.clear(); metadataUpdates.clear(); authenticationErrors.clear(); + totalRequestCount.set(0); } public boolean hasPendingMetadataUpdates() { @@ -462,7 +463,7 @@ public ClientRequest newClientRequest(String nodeId, AbstractRequest.Builder @Override public ClientRequest newClientRequest(String nodeId, AbstractRequest.Builder requestBuilder, long createdTimeMs, boolean expectResponse, RequestCompletionHandler callback) { - requestCount.incrementAndGet(); + totalRequestCount.incrementAndGet(); return new ClientRequest(nodeId, requestBuilder, 0, "mockClientId", createdTimeMs, expectResponse, callback); } @@ -508,6 +509,6 @@ private static class MetadataUpdate { // visible for testing public int requestCount() { - return requestCount.get(); + return totalRequestCount.get(); } } From bfdc85e1b9398d4670a774ff77812f23cece88c0 Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Thu, 5 Apr 2018 00:02:10 -0700 Subject: [PATCH 4/4] Rename requestCount to totalRequestCount --- .../src/test/java/org/apache/kafka/clients/MockClient.java | 2 +- .../apache/kafka/clients/admin/KafkaAdminClientTest.java | 6 +++--- 2 files changed, 4 insertions(+), 4 deletions(-) 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 f3123dc43181d..a73175c995430 100644 --- a/clients/src/test/java/org/apache/kafka/clients/MockClient.java +++ b/clients/src/test/java/org/apache/kafka/clients/MockClient.java @@ -508,7 +508,7 @@ private static class MetadataUpdate { } // visible for testing - public int requestCount() { + 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 f382ea96ec163..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,14 +230,14 @@ public void testInvalidTopicNames() throws Exception { for (String sillyTopicName : sillyTopicNames) { assertFutureError(deleteFutures.get(sillyTopicName), InvalidTopicException.class); } - assertEquals(0, env.kafkaClient().requestCount()); + 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().requestCount()); + assertEquals(0, env.kafkaClient().totalRequestCount()); List newTopics = new ArrayList<>(); for (String sillyTopicName : sillyTopicNames) { @@ -249,7 +249,7 @@ public void testInvalidTopicNames() throws Exception { for (String sillyTopicName : sillyTopicNames) { assertFutureError(createFutures .get(sillyTopicName), InvalidTopicException.class); } - assertEquals(0, env.kafkaClient().requestCount()); + assertEquals(0, env.kafkaClient().totalRequestCount()); } }