From e610d1c9d2df9c71b9867d62440b1f9a1d4bc866 Mon Sep 17 00:00:00 2001 From: Dmitry Werner Date: Sun, 3 Mar 2024 01:51:12 +0500 Subject: [PATCH 1/7] KAFKA-16246: Cleanups in ConsoleConsumer Removed Optional where it makes sense and removed the argument checking logic that duplicates existing in ConsoleConsumerOptions#checkRequiredArgs(). --- .../kafka/tools/consumer/ConsoleConsumer.java | 53 ++++------- .../tools/consumer/ConsoleConsumerTest.java | 89 ++++++++++--------- 2 files changed, 65 insertions(+), 77 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumer.java b/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumer.java index f84fb88c23f5f..d663b2964d52e 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumer.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumer.java @@ -22,8 +22,6 @@ import java.util.Iterator; import java.util.Map; import java.util.Optional; -import java.util.OptionalInt; -import java.util.OptionalLong; import java.util.concurrent.CountDownLatch; import java.util.regex.Pattern; import java.util.Collections; @@ -68,11 +66,8 @@ public static void main(String[] args) throws Exception { public static void run(ConsoleConsumerOptions opts) { messageCount = 0; - long timeoutMs = opts.timeoutMs() >= 0 ? opts.timeoutMs() : Long.MAX_VALUE; Consumer consumer = new KafkaConsumer<>(opts.consumerProps(), new ByteArrayDeserializer(), new ByteArrayDeserializer()); - ConsumerWrapper consumerWrapper = opts.partitionArg().isPresent() - ? new ConsumerWrapper(Optional.of(opts.topicArg()), opts.partitionArg(), OptionalLong.of(opts.offsetArg()), Optional.empty(), consumer, timeoutMs) - : new ConsumerWrapper(Optional.of(opts.topicArg()), OptionalInt.empty(), OptionalLong.empty(), Optional.ofNullable(opts.includedTopicsArg()), consumer, timeoutMs); + ConsumerWrapper consumerWrapper = new ConsumerWrapper(opts, consumer); addShutdownHook(consumerWrapper, opts); @@ -148,43 +143,27 @@ static boolean checkErr(PrintStream output) { } public static class ConsumerWrapper { - final Optional topic; - final OptionalInt partitionId; - final OptionalLong offset; - final Optional includedTopics; - final Consumer consumer; - final long timeoutMs; final Time time = Time.SYSTEM; + final long timeoutMs; + final Consumer consumer; Iterator> recordIter = Collections.emptyIterator(); - public ConsumerWrapper(Optional topic, - OptionalInt partitionId, - OptionalLong offset, - Optional includedTopics, - Consumer consumer, - long timeoutMs) { - this.topic = topic; - this.partitionId = partitionId; - this.offset = offset; - this.includedTopics = includedTopics; + public ConsumerWrapper(ConsoleConsumerOptions opts, Consumer consumer) { + boolean isPartitionArgPresented = opts.partitionArg().isPresent(); + Optional topic = Optional.ofNullable(opts.topicArg()); + Optional includedTopics = isPartitionArgPresented ? Optional.empty() : Optional.ofNullable(opts.includedTopicsArg()); this.consumer = consumer; - this.timeoutMs = timeoutMs; - - if (topic.isPresent() && partitionId.isPresent() && offset.isPresent() && !includedTopics.isPresent()) { - seek(topic.get(), partitionId.getAsInt(), offset.getAsLong()); - } else if (topic.isPresent() && partitionId.isPresent() && !offset.isPresent() && !includedTopics.isPresent()) { - // default to latest if no offset is provided - seek(topic.get(), partitionId.getAsInt(), ListOffsetsRequest.LATEST_TIMESTAMP); - } else if (topic.isPresent() && !partitionId.isPresent() && !offset.isPresent() && !includedTopics.isPresent()) { - consumer.subscribe(Collections.singletonList(topic.get())); - } else if (!topic.isPresent() && !partitionId.isPresent() && !offset.isPresent() && includedTopics.isPresent()) { - consumer.subscribe(Pattern.compile(includedTopics.get())); + timeoutMs = opts.timeoutMs() >= 0 ? opts.timeoutMs() : Long.MAX_VALUE; + + if (topic.isPresent()) { + if (isPartitionArgPresented) { + seek(topic.get(), opts.partitionArg().getAsInt(), opts.offsetArg()); + } else { + consumer.subscribe(Collections.singletonList(topic.get())); + } } else { - throw new IllegalArgumentException("An invalid combination of arguments is provided. " + - "Exactly one of 'topic' or 'include' must be provided. " + - "If 'topic' is provided, an optional 'partition' may also be provided. " + - "If 'partition' is provided, an optional 'offset' may also be provided, otherwise, consumption starts from the end of the partition."); + includedTopics.ifPresent(topics -> consumer.subscribe(Pattern.compile(topics))); } } diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java index 008893f9c505a..9df8fbd88abbe 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java @@ -24,21 +24,18 @@ import org.apache.kafka.common.MessageFormatter; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.TimeoutException; -import org.apache.kafka.common.requests.ListOffsetsRequest; import org.apache.kafka.common.utils.Time; import org.apache.kafka.server.util.MockTime; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import java.io.IOException; import java.io.PrintStream; import java.time.Duration; import java.util.Arrays; import java.util.Collections; import java.util.HashMap; import java.util.Map; -import java.util.Optional; -import java.util.OptionalInt; -import java.util.OptionalLong; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; @@ -58,8 +55,7 @@ public void setup() { } @Test - public void shouldThrowTimeoutExceptionWhenTimeoutIsReached() { - String topic = "test"; + public void shouldThrowTimeoutExceptionWhenTimeoutIsReached() throws IOException { final Time time = new MockTime(); final int timeoutMs = 1000; @@ -71,20 +67,22 @@ public void shouldThrowTimeoutExceptionWhenTimeoutIsReached() { return ConsumerRecords.EMPTY; }); + String[] args = new String[]{ + "--bootstrap-server", "localhost:9092", + "--topic", "test", + "--timeout-ms", String.valueOf(timeoutMs) + }; + ConsoleConsumer.ConsumerWrapper consumer = new ConsoleConsumer.ConsumerWrapper( - Optional.of(topic), - OptionalInt.empty(), - OptionalLong.empty(), - Optional.empty(), - mockConsumer, - timeoutMs + new ConsoleConsumerOptions(args), + mockConsumer ); assertThrows(TimeoutException.class, consumer::receive); } @Test - public void shouldResetUnConsumedOffsetsBeforeExit() { + public void shouldResetUnConsumedOffsetsBeforeExit() throws IOException { String topic = "test"; int maxMessages = 123; int totalMessages = 700; @@ -94,13 +92,16 @@ public void shouldResetUnConsumedOffsetsBeforeExit() { TopicPartition tp1 = new TopicPartition(topic, 0); TopicPartition tp2 = new TopicPartition(topic, 1); + String[] args = new String[]{ + "--bootstrap-server", "localhost:9092", + "--topic", topic, + "--timeout-ms", "1000" + }; + ConsoleConsumer.ConsumerWrapper consumer = new ConsoleConsumer.ConsumerWrapper( - Optional.of(topic), - OptionalInt.empty(), - OptionalLong.empty(), - Optional.empty(), - mockConsumer, - 1000L); + new ConsoleConsumerOptions(args), + mockConsumer + ); mockConsumer.rebalance(Arrays.asList(tp1, tp2)); Map offsets = new HashMap<>(); @@ -165,43 +166,51 @@ public void shouldStopWhenOutputCheckErrorFails() { @Test @SuppressWarnings("unchecked") - public void shouldSeekWhenOffsetIsSet() { + public void shouldSeekWhenOffsetIsSet() throws IOException { Consumer mockConsumer = mock(Consumer.class); TopicPartition tp0 = new TopicPartition("test", 0); + String[] args = new String[]{ + "--bootstrap-server", "localhost:9092", + "--topic", tp0.topic(), + "--partition", String.valueOf(tp0.partition()), + "--timeout-ms", "1000" + }; + ConsoleConsumer.ConsumerWrapper consumer = new ConsoleConsumer.ConsumerWrapper( - Optional.of(tp0.topic()), - OptionalInt.of(tp0.partition()), - OptionalLong.empty(), - Optional.empty(), - mockConsumer, - 1000L); + new ConsoleConsumerOptions(args), + mockConsumer + ); verify(mockConsumer).assign(eq(Collections.singletonList(tp0))); verify(mockConsumer).seekToEnd(eq(Collections.singletonList(tp0))); consumer.cleanup(); reset(mockConsumer); - consumer = new ConsoleConsumer.ConsumerWrapper( - Optional.of(tp0.topic()), - OptionalInt.of(tp0.partition()), - OptionalLong.of(123L), - Optional.empty(), - mockConsumer, - 1000L); + args = new String[]{ + "--bootstrap-server", "localhost:9092", + "--topic", tp0.topic(), + "--partition", String.valueOf(tp0.partition()), + "--offset", "123", + "--timeout-ms", "1000" + }; + + consumer = new ConsoleConsumer.ConsumerWrapper(new ConsoleConsumerOptions(args), mockConsumer); verify(mockConsumer).assign(eq(Collections.singletonList(tp0))); verify(mockConsumer).seek(eq(tp0), eq(123L)); consumer.cleanup(); reset(mockConsumer); - consumer = new ConsoleConsumer.ConsumerWrapper( - Optional.of(tp0.topic()), - OptionalInt.of(tp0.partition()), - OptionalLong.of(ListOffsetsRequest.EARLIEST_TIMESTAMP), - Optional.empty(), - mockConsumer, - 1000L); + args = new String[]{ + "--bootstrap-server", "localhost:9092", + "--topic", tp0.topic(), + "--partition", String.valueOf(tp0.partition()), + "--offset", "earliest", + "--timeout-ms", "1000" + }; + + consumer = new ConsoleConsumer.ConsumerWrapper(new ConsoleConsumerOptions(args), mockConsumer); verify(mockConsumer).assign(eq(Collections.singletonList(tp0))); verify(mockConsumer).seekToBeginning(eq(Collections.singletonList(tp0))); From 2f96fe816326315e0d1da47a0159a0664243f451 Mon Sep 17 00:00:00 2001 From: Dmitry Werner Date: Sun, 3 Mar 2024 04:03:06 +0500 Subject: [PATCH 2/7] KAFKA-16246: Cleanups in ConsoleConsumer added test for NPE fix. --- .../tools/consumer/ConsoleConsumerTest.java | 32 +++++++++++++++---- 1 file changed, 26 insertions(+), 6 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java index 9df8fbd88abbe..f67c1f582aaf3 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java @@ -36,6 +36,7 @@ import java.util.Collections; import java.util.HashMap; import java.util.Map; +import java.util.regex.Pattern; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; @@ -74,8 +75,8 @@ public void shouldThrowTimeoutExceptionWhenTimeoutIsReached() throws IOException }; ConsoleConsumer.ConsumerWrapper consumer = new ConsoleConsumer.ConsumerWrapper( - new ConsoleConsumerOptions(args), - mockConsumer + new ConsoleConsumerOptions(args), + mockConsumer ); assertThrows(TimeoutException.class, consumer::receive); @@ -99,8 +100,8 @@ public void shouldResetUnConsumedOffsetsBeforeExit() throws IOException { }; ConsoleConsumer.ConsumerWrapper consumer = new ConsoleConsumer.ConsumerWrapper( - new ConsoleConsumerOptions(args), - mockConsumer + new ConsoleConsumerOptions(args), + mockConsumer ); mockConsumer.rebalance(Arrays.asList(tp1, tp2)); @@ -178,8 +179,8 @@ public void shouldSeekWhenOffsetIsSet() throws IOException { }; ConsoleConsumer.ConsumerWrapper consumer = new ConsoleConsumer.ConsumerWrapper( - new ConsoleConsumerOptions(args), - mockConsumer + new ConsoleConsumerOptions(args), + mockConsumer ); verify(mockConsumer).assign(eq(Collections.singletonList(tp0))); @@ -217,4 +218,23 @@ public void shouldSeekWhenOffsetIsSet() throws IOException { consumer.cleanup(); reset(mockConsumer); } + + @Test + @SuppressWarnings("unchecked") + public void shouldWorkWithoutTopicOption() throws IOException { + Consumer mockConsumer = mock(Consumer.class); + + String[] args = new String[]{ + "--bootstrap-server", "localhost:9092", + "--include", "includeTest*", + "--from-beginning" + }; + + new ConsoleConsumer.ConsumerWrapper( + new ConsoleConsumerOptions(args), + mockConsumer + ); + + verify(mockConsumer).subscribe(any(Pattern.class)); + } } From be28d6a027404af2f733f81e1c192526cdf60991 Mon Sep 17 00:00:00 2001 From: Dmitry Werner Date: Sun, 3 Mar 2024 04:09:40 +0500 Subject: [PATCH 3/7] KAFKA-16246: Cleanups in ConsoleConsumer minor --- .../org/apache/kafka/tools/consumer/ConsoleConsumerTest.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java index f67c1f582aaf3..ab0710d270506 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java @@ -230,10 +230,7 @@ public void shouldWorkWithoutTopicOption() throws IOException { "--from-beginning" }; - new ConsoleConsumer.ConsumerWrapper( - new ConsoleConsumerOptions(args), - mockConsumer - ); + new ConsoleConsumer.ConsumerWrapper(new ConsoleConsumerOptions(args), mockConsumer); verify(mockConsumer).subscribe(any(Pattern.class)); } From 201316e82840575c8b58f4ea6429f2a0d5c8fff5 Mon Sep 17 00:00:00 2001 From: Dmitry Werner Date: Sun, 3 Mar 2024 04:12:27 +0500 Subject: [PATCH 4/7] KAFKA-16246: Cleanups in ConsoleConsumer minor --- .../apache/kafka/tools/consumer/ConsoleConsumerTest.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java index ab0710d270506..f4fa6ac3be221 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerTest.java @@ -230,8 +230,12 @@ public void shouldWorkWithoutTopicOption() throws IOException { "--from-beginning" }; - new ConsoleConsumer.ConsumerWrapper(new ConsoleConsumerOptions(args), mockConsumer); + ConsoleConsumer.ConsumerWrapper consumer = new ConsoleConsumer.ConsumerWrapper( + new ConsoleConsumerOptions(args), + mockConsumer + ); verify(mockConsumer).subscribe(any(Pattern.class)); + consumer.cleanup(); } } From 32e50b801593dd619b8b4071bd823859da6e988d Mon Sep 17 00:00:00 2001 From: Dmitry Werner Date: Sun, 3 Mar 2024 11:08:57 +0500 Subject: [PATCH 5/7] KAFKA-16246: Cleanups in ConsoleConsumer minor --- .../org/apache/kafka/tools/consumer/ConsoleConsumer.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumer.java b/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumer.java index d663b2964d52e..70ed1515d6d5a 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumer.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumer.java @@ -150,14 +150,13 @@ public static class ConsumerWrapper { Iterator> recordIter = Collections.emptyIterator(); public ConsumerWrapper(ConsoleConsumerOptions opts, Consumer consumer) { - boolean isPartitionArgPresented = opts.partitionArg().isPresent(); Optional topic = Optional.ofNullable(opts.topicArg()); - Optional includedTopics = isPartitionArgPresented ? Optional.empty() : Optional.ofNullable(opts.includedTopicsArg()); + Optional includedTopics = Optional.ofNullable(opts.includedTopicsArg()); this.consumer = consumer; timeoutMs = opts.timeoutMs() >= 0 ? opts.timeoutMs() : Long.MAX_VALUE; if (topic.isPresent()) { - if (isPartitionArgPresented) { + if (opts.partitionArg().isPresent()) { seek(topic.get(), opts.partitionArg().getAsInt(), opts.offsetArg()); } else { consumer.subscribe(Collections.singletonList(topic.get())); From a63b0a4b00463a63034a1d53436589d651d72155 Mon Sep 17 00:00:00 2001 From: Dmitry Werner Date: Wed, 6 Mar 2024 23:12:29 +0500 Subject: [PATCH 6/7] KAFKA-16246: Cleanups in ConsoleConsumer fix review comments --- .../kafka/tools/consumer/ConsoleConsumer.java | 7 ++- .../consumer/ConsoleConsumerOptions.java | 49 ++++++++++++----- .../consumer/ConsoleConsumerOptionsTest.java | 54 ++++++++++++++----- 3 files changed, 80 insertions(+), 30 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumer.java b/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumer.java index 70ed1515d6d5a..bb5ab1443ed49 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumer.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumer.java @@ -150,10 +150,9 @@ public static class ConsumerWrapper { Iterator> recordIter = Collections.emptyIterator(); public ConsumerWrapper(ConsoleConsumerOptions opts, Consumer consumer) { - Optional topic = Optional.ofNullable(opts.topicArg()); - Optional includedTopics = Optional.ofNullable(opts.includedTopicsArg()); this.consumer = consumer; - timeoutMs = opts.timeoutMs() >= 0 ? opts.timeoutMs() : Long.MAX_VALUE; + timeoutMs = opts.timeoutMs(); + Optional topic = opts.topicArg(); if (topic.isPresent()) { if (opts.partitionArg().isPresent()) { @@ -162,7 +161,7 @@ public ConsumerWrapper(ConsoleConsumerOptions opts, Consumer con consumer.subscribe(Collections.singletonList(topic.get())); } } else { - includedTopics.ifPresent(topics -> consumer.subscribe(Pattern.compile(topics))); + opts.includedTopicsArg().ifPresent(topics -> consumer.subscribe(Pattern.compile(topics))); } } diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptions.java b/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptions.java index a713afb2bf22c..64c74d7bc4ab9 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptions.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptions.java @@ -34,7 +34,7 @@ import java.util.List; import java.util.Locale; import java.util.Map; -import java.util.Objects; +import java.util.Optional; import java.util.OptionalInt; import java.util.Properties; import java.util.Random; @@ -55,7 +55,7 @@ public final class ConsoleConsumerOptions extends CommandDefaultOptions { private final OptionSpec messageFormatterConfigOpt; private final OptionSpec resetBeginningOpt; private final OptionSpec maxMessagesOpt; - private final OptionSpec timeoutMsOpt; + private final OptionSpec timeoutMsOpt; private final OptionSpec skipMessageOnErrorOpt; private final OptionSpec bootstrapServerOpt; private final OptionSpec keyDeserializerOpt; @@ -66,6 +66,7 @@ public final class ConsoleConsumerOptions extends CommandDefaultOptions { private final Properties consumerProps; private final long offset; + private final long timeoutMs; private final MessageFormatter formatter; public ConsoleConsumerOptions(String[] args) throws IOException { @@ -139,7 +140,7 @@ public ConsoleConsumerOptions(String[] args) throws IOException { timeoutMsOpt = parser.accepts("timeout-ms", "If specified, exit if no message is available for consumption for the specified interval.") .withRequiredArg() .describedAs("timeout_ms") - .ofType(Integer.class); + .ofType(Long.class); skipMessageOnErrorOpt = parser.accepts("skip-message-on-error", "If there is an error when processing a message, " + "skip it instead of halt."); bootstrapServerOpt = parser.accepts("bootstrap-server", "REQUIRED: The server(s) to connect to.") @@ -184,12 +185,13 @@ public ConsoleConsumerOptions(String[] args) throws IOException { Set groupIdsProvided = checkConsumerGroup(consumerPropsFromFile, extraConsumerProps); consumerProps = buildConsumerProps(consumerPropsFromFile, extraConsumerProps, groupIdsProvided); offset = parseOffset(); + timeoutMs = parseTimeoutMs(); formatter = buildFormatter(); } private void checkRequiredArgs() { - List topicOrFilterArgs = new ArrayList<>(Arrays.asList(topicArg(), includedTopicsArg())); - topicOrFilterArgs.removeIf(Objects::isNull); + List> topicOrFilterArgs = new ArrayList<>(Arrays.asList(topicArg(), includedTopicsArg())); + topicOrFilterArgs.removeIf(arg -> !arg.isPresent()); // user need to specify value for either --topic or one of the include filters options (--include or --whitelist) if (topicOrFilterArgs.size() != 1) { CommandLineUtils.printUsageAndExit(parser, "Exactly one of --include/--topic is required. " + @@ -322,6 +324,20 @@ private void invalidOffset(String offset) { "'earliest', 'latest', or a non-negative long."); } + private long parseTimeoutMs() { + long timeout; + if (options.has(timeoutMsOpt)) { + timeout = options.valueOf(timeoutMsOpt); + if (timeout < 0) { + CommandLineUtils.printUsageAndExit(parser, "The provided timeout-ms value '" + timeout + + "' is incorrect. Valid value are a non-negative long."); + } + } else { + timeout = Long.MAX_VALUE; + } + return timeout; + } + private MessageFormatter buildFormatter() { MessageFormatter formatter = null; try { @@ -365,16 +381,19 @@ OptionalInt partitionArg() { return OptionalInt.empty(); } - String topicArg() { - return options.valueOf(topicOpt); + Optional topicArg() { + if (options.has(topicOpt)) { + return Optional.of(options.valueOf(topicOpt)); + } + return Optional.empty(); } int maxMessages() { return options.has(maxMessagesOpt) ? options.valueOf(maxMessagesOpt) : -1; } - int timeoutMs() { - return options.has(timeoutMsOpt) ? options.valueOf(timeoutMsOpt) : -1; + long timeoutMs() { + return timeoutMs; } boolean enableSystestEventsLogging() { @@ -385,10 +404,14 @@ String bootstrapServer() { return options.valueOf(bootstrapServerOpt); } - String includedTopicsArg() { - return options.has(includeOpt) - ? options.valueOf(includeOpt) - : options.valueOf(whitelistOpt); + Optional includedTopicsArg() { + if (options.has(includeOpt)) { + return Optional.of(options.valueOf(includeOpt)); + } else if (options.has(whitelistOpt)) { + return Optional.of(options.valueOf(whitelistOpt)); + } else { + return Optional.empty(); + } } Properties formatterArgs() throws IOException { diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptionsTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptionsTest.java index 523122c4cdd98..5fdb90c7aad30 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptionsTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptionsTest.java @@ -48,12 +48,12 @@ public void shouldParseValidConsumerValidConfig() throws IOException { ConsoleConsumerOptions config = new ConsoleConsumerOptions(args); assertEquals("localhost:9092", config.bootstrapServer()); - assertEquals("test", config.topicArg()); + assertEquals("test", config.topicArg().orElse("")); assertTrue(config.fromBeginning()); assertFalse(config.enableSystestEventsLogging()); assertFalse(config.skipMessageOnError()); assertEquals(-1, config.maxMessages()); - assertEquals(-1, config.timeoutMs()); + assertEquals(Long.MAX_VALUE, config.timeoutMs()); } @Test @@ -67,7 +67,7 @@ public void shouldParseIncludeArgument() throws IOException { ConsoleConsumerOptions config = new ConsoleConsumerOptions(args); assertEquals("localhost:9092", config.bootstrapServer()); - assertEquals("includeTest*", config.includedTopicsArg()); + assertEquals("includeTest*", config.includedTopicsArg().orElse("")); assertTrue(config.fromBeginning()); } @@ -82,7 +82,7 @@ public void shouldParseWhitelistArgument() throws IOException { ConsoleConsumerOptions config = new ConsoleConsumerOptions(args); assertEquals("localhost:9092", config.bootstrapServer()); - assertEquals("whitelistTest*", config.includedTopicsArg()); + assertEquals("whitelistTest*", config.includedTopicsArg().orElse("")); assertTrue(config.fromBeginning()); } @@ -96,7 +96,7 @@ public void shouldIgnoreWhitelistArgumentIfIncludeSpecified() throws IOException }; ConsoleConsumerOptions config = new ConsoleConsumerOptions(args); assertEquals("localhost:9092", config.bootstrapServer()); - assertEquals("includeTest*", config.includedTopicsArg()); + assertEquals("includeTest*", config.includedTopicsArg().orElse("")); assertTrue(config.fromBeginning()); } @@ -112,7 +112,7 @@ public void shouldParseValidSimpleConsumerValidConfigWithNumericOffset() throws ConsoleConsumerOptions config = new ConsoleConsumerOptions(args); assertEquals("localhost:9092", config.bootstrapServer()); - assertEquals("test", config.topicArg()); + assertEquals("test", config.topicArg().orElse("")); assertTrue(config.partitionArg().isPresent()); assertEquals(0, config.partitionArg().getAsInt()); assertEquals(3, config.offsetArg()); @@ -191,7 +191,7 @@ public void shouldParseValidSimpleConsumerValidConfigWithStringOffset() throws E ConsoleConsumerOptions config = new ConsoleConsumerOptions(args); assertEquals("localhost:9092", config.bootstrapServer()); - assertEquals("test", config.topicArg()); + assertEquals("test", config.topicArg().orElse("")); assertTrue(config.partitionArg().isPresent()); assertEquals(0, config.partitionArg().getAsInt()); assertEquals(-1, config.offsetArg()); @@ -211,7 +211,7 @@ public void shouldParseValidConsumerConfigWithAutoOffsetResetLatest() throws IOE Properties consumerProperties = config.consumerProps(); assertEquals("localhost:9092", config.bootstrapServer()); - assertEquals("test", config.topicArg()); + assertEquals("test", config.topicArg().orElse("")); assertFalse(config.fromBeginning()); assertEquals("latest", consumerProperties.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)); } @@ -228,7 +228,7 @@ public void shouldParseValidConsumerConfigWithAutoOffsetResetEarliest() throws I Properties consumerProperties = config.consumerProps(); assertEquals("localhost:9092", config.bootstrapServer()); - assertEquals("test", config.topicArg()); + assertEquals("test", config.topicArg().orElse("")); assertFalse(config.fromBeginning()); assertEquals("earliest", consumerProperties.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)); } @@ -246,7 +246,7 @@ public void shouldParseValidConsumerConfigWithAutoOffsetResetAndMatchingFromBegi Properties consumerProperties = config.consumerProps(); assertEquals("localhost:9092", config.bootstrapServer()); - assertEquals("test", config.topicArg()); + assertEquals("test", config.topicArg().orElse("")); assertTrue(config.fromBeginning()); assertEquals("earliest", consumerProperties.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)); } @@ -262,7 +262,7 @@ public void shouldParseValidConsumerConfigWithNoOffsetReset() throws IOException Properties consumerProperties = config.consumerProps(); assertEquals("localhost:9092", config.bootstrapServer()); - assertEquals("test", config.topicArg()); + assertEquals("test", config.topicArg().orElse("")); assertFalse(config.fromBeginning()); assertEquals("latest", consumerProperties.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)); } @@ -442,7 +442,7 @@ public void shouldParseGroupIdFromBeginningGivenTogether() throws IOException { ConsoleConsumerOptions config = new ConsoleConsumerOptions(args); assertEquals("localhost:9092", config.bootstrapServer()); - assertEquals("test", config.topicArg()); + assertEquals("test", config.topicArg().orElse("")); assertEquals(-2, config.offsetArg()); assertTrue(config.fromBeginning()); @@ -455,7 +455,7 @@ public void shouldParseGroupIdFromBeginningGivenTogether() throws IOException { config = new ConsoleConsumerOptions(args); assertEquals("localhost:9092", config.bootstrapServer()); - assertEquals("test", config.topicArg()); + assertEquals("test", config.topicArg().orElse("")); assertEquals(-1, config.offsetArg()); assertFalse(config.fromBeginning()); } @@ -618,4 +618,32 @@ public void testParseOffset() throws Exception { Exit.resetExitProcedure(); } } + + @Test + public void testParseTimeoutMs() throws Exception { + Exit.setExitProcedure((code, message) -> { + throw new IllegalArgumentException(message); + }); + + try { + String[] negativeTimeoutMs = new String[]{ + "--bootstrap-server", "localhost:9092", + "--topic", "test", + "--partition", "0", + "--timeout-ms", "-1" + }; + assertThrows(IllegalArgumentException.class, () -> new ConsoleConsumerOptions(negativeTimeoutMs)); + + String[] validTimeoutMs = new String[]{ + "--bootstrap-server", "localhost:9092", + "--topic", "test", + "--partition", "0", + "--timeout-ms", "100" + }; + ConsoleConsumerOptions config = new ConsoleConsumerOptions(validTimeoutMs); + assertEquals(100, config.timeoutMs()); + } finally { + Exit.resetExitProcedure(); + } + } } From bdf8f5b0610e3677d0bc140a79f455d71826ca3d Mon Sep 17 00:00:00 2001 From: Dmitry Werner Date: Thu, 7 Mar 2024 00:00:10 +0500 Subject: [PATCH 7/7] KAFKA-16246: Cleanups in ConsoleConsumer fix review comments --- .../consumer/ConsoleConsumerOptions.java | 28 +++---------- .../consumer/ConsoleConsumerOptionsTest.java | 42 +++++++++---------- 2 files changed, 26 insertions(+), 44 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptions.java b/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptions.java index 64c74d7bc4ab9..aa37919515450 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptions.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptions.java @@ -325,17 +325,8 @@ private void invalidOffset(String offset) { } private long parseTimeoutMs() { - long timeout; - if (options.has(timeoutMsOpt)) { - timeout = options.valueOf(timeoutMsOpt); - if (timeout < 0) { - CommandLineUtils.printUsageAndExit(parser, "The provided timeout-ms value '" + timeout + - "' is incorrect. Valid value are a non-negative long."); - } - } else { - timeout = Long.MAX_VALUE; - } - return timeout; + long timeout = options.has(timeoutMsOpt) ? options.valueOf(timeoutMsOpt) : -1; + return timeout >= 0 ? timeout : Long.MAX_VALUE; } private MessageFormatter buildFormatter() { @@ -382,10 +373,7 @@ OptionalInt partitionArg() { } Optional topicArg() { - if (options.has(topicOpt)) { - return Optional.of(options.valueOf(topicOpt)); - } - return Optional.empty(); + return options.has(topicOpt) ? Optional.of(options.valueOf(topicOpt)) : Optional.empty(); } int maxMessages() { @@ -405,13 +393,9 @@ String bootstrapServer() { } Optional includedTopicsArg() { - if (options.has(includeOpt)) { - return Optional.of(options.valueOf(includeOpt)); - } else if (options.has(whitelistOpt)) { - return Optional.of(options.valueOf(whitelistOpt)); - } else { - return Optional.empty(); - } + return options.has(includeOpt) + ? Optional.of(options.valueOf(includeOpt)) + : Optional.ofNullable(options.valueOf(whitelistOpt)); } Properties formatterArgs() throws IOException { diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptionsTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptionsTest.java index 5fdb90c7aad30..3242b642cdb82 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptionsTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/ConsoleConsumerOptionsTest.java @@ -621,29 +621,27 @@ public void testParseOffset() throws Exception { @Test public void testParseTimeoutMs() throws Exception { - Exit.setExitProcedure((code, message) -> { - throw new IllegalArgumentException(message); - }); + String[] withoutTimeoutMs = new String[]{ + "--bootstrap-server", "localhost:9092", + "--topic", "test", + "--partition", "0" + }; + assertEquals(Long.MAX_VALUE, new ConsoleConsumerOptions(withoutTimeoutMs).timeoutMs()); - try { - String[] negativeTimeoutMs = new String[]{ - "--bootstrap-server", "localhost:9092", - "--topic", "test", - "--partition", "0", - "--timeout-ms", "-1" - }; - assertThrows(IllegalArgumentException.class, () -> new ConsoleConsumerOptions(negativeTimeoutMs)); + String[] negativeTimeoutMs = new String[]{ + "--bootstrap-server", "localhost:9092", + "--topic", "test", + "--partition", "0", + "--timeout-ms", "-100" + }; + assertEquals(Long.MAX_VALUE, new ConsoleConsumerOptions(negativeTimeoutMs).timeoutMs()); - String[] validTimeoutMs = new String[]{ - "--bootstrap-server", "localhost:9092", - "--topic", "test", - "--partition", "0", - "--timeout-ms", "100" - }; - ConsoleConsumerOptions config = new ConsoleConsumerOptions(validTimeoutMs); - assertEquals(100, config.timeoutMs()); - } finally { - Exit.resetExitProcedure(); - } + String[] validTimeoutMs = new String[]{ + "--bootstrap-server", "localhost:9092", + "--topic", "test", + "--partition", "0", + "--timeout-ms", "100" + }; + assertEquals(100, new ConsoleConsumerOptions(validTimeoutMs).timeoutMs()); } }