From c990b265a0781dc567e4149a0f6b377053b7192e Mon Sep 17 00:00:00 2001 From: nizhikov Date: Tue, 28 Nov 2023 18:46:20 +0300 Subject: [PATCH 01/18] KAFKA-14589 ConsumerGroupCommand options and case classes rewritten in java --- .../org/apache/kafka/common/utils/Utils.java | 10 + .../ConsumerGroupCommandOptions.java | 263 ++++++++++++++++++ .../kafka/tools/consumergroup/GroupState.java | 35 +++ .../consumergroup/MemberAssignmentState.java | 42 +++ .../PartitionAssignmentState.java | 50 ++++ 5 files changed, 400 insertions(+) create mode 100644 tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java create mode 100644 tools/src/main/java/org/apache/kafka/tools/consumergroup/GroupState.java create mode 100644 tools/src/main/java/org/apache/kafka/tools/consumergroup/MemberAssignmentState.java create mode 100644 tools/src/main/java/org/apache/kafka/tools/consumergroup/PartitionAssignmentState.java diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 6a0913d3c2da1..0a6461910a632 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -593,6 +593,16 @@ public static String join(T[] strs, String separator) { return join(Arrays.asList(strs), separator); } + /** + * Create a string representation of a collection joined by ", ". + * @param collection The list of items + * @return The string representation. + */ + public static String join(Collection collection) { + Objects.requireNonNull(collection); + return mkString(collection.stream(), "", "", ", "); + } + /** * Create a string representation of a collection joined by the given separator * @param collection The list of items diff --git a/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java b/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java new file mode 100644 index 0000000000000..8ab65746e7b06 --- /dev/null +++ b/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java @@ -0,0 +1,263 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.tools.consumergroup; + +import joptsimple.OptionSpec; +import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.server.util.CommandDefaultOptions; +import org.apache.kafka.server.util.CommandLineUtils; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.Arrays; +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +public class ConsumerGroupCommandOptions extends CommandDefaultOptions { + public static final Logger LOGGER = LoggerFactory.getLogger(ConsumerGroupCommandOptions.class); + + public static final String BOOTSTRAP_SERVER_DOC = "REQUIRED: The server(s) to connect to."; + public static final String GROUP_DOC = "The consumer group we wish to act on."; + public static final String TOPIC_DOC = "The topic whose consumer group information should be deleted or topic whose should be included in the reset offset process. " + + "In `reset-offsets` case, partitions can be specified using this format: `topic1:0,1,2`, where 0,1,2 are the partition to be included in the process. " + + "Reset-offsets also supports multiple topic inputs."; + public static final String ALL_TOPICS_DOC = "Consider all topics assigned to a group in the `reset-offsets` process."; + public static final String LIST_DOC = "List all consumer groups."; + public static final String DESCRIBE_DOC = "Describe consumer group and list offset lag (number of messages not yet processed) related to given group."; + public static final String ALL_GROUPS_DOC = "Apply to all consumer groups."; + public static final String NL = System.getProperty("line.separator"); + public static final String DELETE_DOC = "Pass in groups to delete topic partition offsets and ownership information " + + "over the entire consumer group. For instance --group g1 --group g2"; + public static final String TIMEOUT_MS_DOC = "The timeout that can be set for some use cases. For example, it can be used when describing the group " + + "to specify the maximum amount of time in milliseconds to wait before the group stabilizes (when the group is just created, " + + "or is going through some changes)."; + public static final String COMMAND_CONFIG_DOC = "Property file containing configs to be passed to Admin Client and Consumer."; + public static final String RESET_OFFSETS_DOC = "Reset offsets of consumer group. Supports one consumer group at the time, and instances should be inactive" + NL + + "Has 2 execution options: --dry-run (the default) to plan which offsets to reset, and --execute to update the offsets. " + + "Additionally, the --export option is used to export the results to a CSV format." + NL + + "You must choose one of the following reset specifications: --to-datetime, --by-duration, --to-earliest, " + + "--to-latest, --shift-by, --from-file, --to-current, --to-offset." + NL + + "To define the scope use --all-topics or --topic. One scope must be specified unless you use '--from-file'."; + public static final String DRY_RUN_DOC = "Only show results without executing changes on Consumer Groups. Supported operations: reset-offsets."; + public static final String EXECUTE_DOC = "Execute operation. Supported operations: reset-offsets."; + public static final String EXPORT_DOC = "Export operation execution to a CSV file. Supported operations: reset-offsets."; + public static final String RESET_TO_OFFSET_DOC = "Reset offsets to a specific offset."; + public static final String RESET_FROM_FILE_DOC = "Reset offsets to values defined in CSV file."; + public static final String RESET_TO_DATETIME_DOC = "Reset offsets to offset from datetime. Format: 'YYYY-MM-DDTHH:mm:SS.sss'"; + public static final String RESET_BY_DURATION_DOC = "Reset offsets to offset by duration from current timestamp. Format: 'PnDTnHnMnS'"; + public static final String RESET_TO_EARLIEST_DOC = "Reset offsets to earliest offset."; + public static final String RESET_TO_LATEST_DOC = "Reset offsets to latest offset."; + public static final String RESET_TO_CURRENT_DOC = "Reset offsets to current offset."; + public static final String RESET_SHIFT_BY_DOC = "Reset offsets shifting current offset by 'n', where 'n' can be positive or negative."; + public static final String MEMBERS_DOC = "Describe members of the group. This option may be used with '--describe' and '--bootstrap-server' options only." + NL + + "Example: --bootstrap-server localhost:9092 --describe --group group1 --members"; + public static final String VERBOSE_DOC = "Provide additional information, if any, when describing the group. This option may be used " + + "with '--offsets'/'--members'/'--state' and '--bootstrap-server' options only." + NL + "Example: --bootstrap-server localhost:9092 --describe --group group1 --members --verbose"; + public static final String OFFSETS_DOC = "Describe the group and list all topic partitions in the group along with their offset lag. " + + "This is the default sub-action of and may be used with '--describe' and '--bootstrap-server' options only." + NL + + "Example: --bootstrap-server localhost:9092 --describe --group group1 --offsets"; + public static final String STATE_DOC = "When specified with '--describe', includes the state of the group." + NL + + "Example: --bootstrap-server localhost:9092 --describe --group group1 --state" + NL + + "When specified with '--list', it displays the state of all groups. It can also be used to list groups with specific states." + NL + + "Example: --bootstrap-server localhost:9092 --list --state stable,empty" + NL + + "This option may be used with '--describe', '--list' and '--bootstrap-server' options only."; + public static final String DELETE_OFFSETS_DOC = "Delete offsets of consumer group. Supports one consumer group at the time, and multiple topics."; + + public final OptionSpec bootstrapServerOpt; + public final OptionSpec groupOpt; + public final OptionSpec topicOpt; + public final OptionSpec allTopicsOpt; + public final OptionSpec listOpt; + public final OptionSpec describeOpt; + public final OptionSpec allGroupsOpt; + public final OptionSpec deleteOpt; + public final OptionSpec timeoutMsOpt; + public final OptionSpec commandConfigOpt; + public final OptionSpec resetOffsetsOpt; + public final OptionSpec deleteOffsetsOpt; + public final OptionSpec dryRunOpt; + public final OptionSpec executeOpt; + public final OptionSpec exportOpt; + public final OptionSpec resetToOffsetOpt; + public final OptionSpec resetFromFileOpt; + public final OptionSpec resetToDatetimeOpt; + public final OptionSpec resetByDurationOpt; + public final OptionSpec resetToEarliestOpt; + public final OptionSpec resetToLatestOpt; + public final OptionSpec resetToCurrentOpt; + public final OptionSpec resetShiftByOpt; + public final OptionSpec membersOpt; + public final OptionSpec verboseOpt; + public final OptionSpec offsetsOpt; + public final OptionSpec stateOpt; + + public final Set> allGroupSelectionScopeOpts; + public final Set> allConsumerGroupLevelOpts; + public final Set> allResetOffsetScenarioOpts; + public final Set> allDeleteOffsetsOpts; + + public ConsumerGroupCommandOptions(String[] args) { + super(args); + + bootstrapServerOpt = parser.accepts("bootstrap-server", BOOTSTRAP_SERVER_DOC) + .withRequiredArg() + .describedAs("server to connect to") + .ofType(String.class); + groupOpt = parser.accepts("group", GROUP_DOC) + .withRequiredArg() + .describedAs("consumer group") + .ofType(String.class); + topicOpt = parser.accepts("topic", TOPIC_DOC) + .withRequiredArg() + .describedAs("topic") + .ofType(String.class); + allTopicsOpt = parser.accepts("all-topics", ALL_TOPICS_DOC); + listOpt = parser.accepts("list", LIST_DOC); + describeOpt = parser.accepts("describe", DESCRIBE_DOC); + allGroupsOpt = parser.accepts("all-groups", ALL_GROUPS_DOC); + deleteOpt = parser.accepts("delete", DELETE_DOC); + timeoutMsOpt = parser.accepts("timeout", TIMEOUT_MS_DOC) + .withRequiredArg() + .describedAs("timeout (ms)") + .ofType(Long.class) + .defaultsTo(5000L); + commandConfigOpt = parser.accepts("command-config", COMMAND_CONFIG_DOC) + .withRequiredArg() + .describedAs("command config property file") + .ofType(String.class); + resetOffsetsOpt = parser.accepts("reset-offsets", RESET_OFFSETS_DOC); + deleteOffsetsOpt = parser.accepts("delete-offsets", DELETE_OFFSETS_DOC); + dryRunOpt = parser.accepts("dry-run", DRY_RUN_DOC); + executeOpt = parser.accepts("execute", EXECUTE_DOC); + exportOpt = parser.accepts("export", EXPORT_DOC); + resetToOffsetOpt = parser.accepts("to-offset", RESET_TO_OFFSET_DOC) + .withRequiredArg() + .describedAs("offset") + .ofType(Long.class); + resetFromFileOpt = parser.accepts("from-file", RESET_FROM_FILE_DOC) + .withRequiredArg() + .describedAs("path to CSV file") + .ofType(String.class); + resetToDatetimeOpt = parser.accepts("to-datetime", RESET_TO_DATETIME_DOC) + .withRequiredArg() + .describedAs("datetime") + .ofType(String.class); + resetByDurationOpt = parser.accepts("by-duration", RESET_BY_DURATION_DOC) + .withRequiredArg() + .describedAs("duration") + .ofType(String.class); + resetToEarliestOpt = parser.accepts("to-earliest", RESET_TO_EARLIEST_DOC); + resetToLatestOpt = parser.accepts("to-latest", RESET_TO_LATEST_DOC); + resetToCurrentOpt = parser.accepts("to-current", RESET_TO_CURRENT_DOC); + resetShiftByOpt = parser.accepts("shift-by", RESET_SHIFT_BY_DOC) + .withRequiredArg() + .describedAs("number-of-offsets") + .ofType(Long.class); + membersOpt = parser.accepts("members", MEMBERS_DOC) + .availableIf(describeOpt); + verboseOpt = parser.accepts("verbose", VERBOSE_DOC) + .availableIf(describeOpt); + offsetsOpt = parser.accepts("offsets", OFFSETS_DOC) + .availableIf(describeOpt); + stateOpt = parser.accepts("state", STATE_DOC) + .availableIf(describeOpt, listOpt) + .withOptionalArg() + .ofType(String.class); + + allGroupSelectionScopeOpts = new HashSet<>(Arrays.asList(groupOpt, allGroupsOpt)); + allConsumerGroupLevelOpts = new HashSet<>(Arrays.asList(listOpt, describeOpt, deleteOpt, resetOffsetsOpt)); + allResetOffsetScenarioOpts = new HashSet<>(Arrays.asList(resetToOffsetOpt, resetShiftByOpt, + resetToDatetimeOpt, resetByDurationOpt, resetToEarliestOpt, resetToLatestOpt, resetToCurrentOpt, resetFromFileOpt)); + allDeleteOffsetsOpts = new HashSet<>(Arrays.asList(groupOpt, topicOpt)); + + options = parser.parse(args); + } + + @SuppressWarnings({"CyclomaticComplexity", "NPathComplexity"}) + public void checkArgs() { + CommandLineUtils.checkRequiredArgs(parser, options, bootstrapServerOpt); + + if (options.has(describeOpt)) { + if (!options.has(groupOpt) && !options.has(allGroupsOpt)) + CommandLineUtils.printUsageAndExit(parser, + "Option $describeOpt takes one of these options: " + Utils.join(allGroupSelectionScopeOpts)); + List> mutuallyExclusiveOpts = Arrays.asList(membersOpt, offsetsOpt, stateOpt); + if (mutuallyExclusiveOpts.stream().mapToInt(o -> options.has(o) ? 1 : 0).sum() > 1) { + CommandLineUtils.printUsageAndExit(parser, + "Option " + describeOpt + " takes at most one of these options: " + Utils.join(mutuallyExclusiveOpts)); + } + if (options.has(stateOpt) && options.valueOf(stateOpt) != null) + CommandLineUtils.printUsageAndExit(parser, + "Option " + describeOpt + " does not take a value for " + stateOpt); + } else { + if (options.has(timeoutMsOpt)) + LOGGER.debug("Option " + timeoutMsOpt + " is applicable only when " + describeOpt + " is used."); + } + + if (options.has(deleteOpt)) { + if (!options.has(groupOpt) && !options.has(allGroupsOpt)) + CommandLineUtils.printUsageAndExit(parser, + "Option " + deleteOpt + " takes one of these options: " + Utils.join(allGroupSelectionScopeOpts)); + if (options.has(topicOpt)) + CommandLineUtils.printUsageAndExit(parser, "The consumer does not support topic-specific offset " + + "deletion from a consumer group."); + } + + if (options.has(deleteOffsetsOpt)) { + if (!options.has(groupOpt) || !options.has(topicOpt)) + CommandLineUtils.printUsageAndExit(parser, + "Option " + deleteOffsetsOpt + " takes the following options: " + Utils.join(allDeleteOffsetsOpts)); + } + + if (options.has(resetOffsetsOpt)) { + if (options.has(dryRunOpt) && options.has(executeOpt)) + CommandLineUtils.printUsageAndExit(parser, "Option " + resetOffsetsOpt + " only accepts one of " + executeOpt + " and " + dryRunOpt); + + if (!options.has(dryRunOpt) && !options.has(executeOpt)) { + System.err.println("WARN: No action will be performed as the --execute option is missing." + + "In a future major release, the default behavior of this command will be to prompt the user before " + + "executing the reset rather than doing a dry run. You should add the --dry-run option explicitly " + + "if you are scripting this command and want to keep the current default behavior without prompting."); + } + + if (!options.has(groupOpt) && !options.has(allGroupsOpt)) + CommandLineUtils.printUsageAndExit(parser, + "Option " + resetOffsetsOpt + " takes one of these options: " + Utils.join(allGroupSelectionScopeOpts)); + CommandLineUtils.checkInvalidArgs(parser, options, resetToOffsetOpt, minus(allResetOffsetScenarioOpts, resetToOffsetOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, resetToDatetimeOpt, minus(allResetOffsetScenarioOpts, resetToDatetimeOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, resetByDurationOpt, minus(allResetOffsetScenarioOpts, resetByDurationOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, resetToEarliestOpt, minus(allResetOffsetScenarioOpts, resetToEarliestOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, resetToLatestOpt, minus(allResetOffsetScenarioOpts, resetToLatestOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, resetToCurrentOpt, minus(allResetOffsetScenarioOpts, resetToCurrentOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, resetShiftByOpt, minus(allResetOffsetScenarioOpts, resetShiftByOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, resetFromFileOpt, minus(allResetOffsetScenarioOpts, resetFromFileOpt)); + } + + CommandLineUtils.checkInvalidArgs(parser, options, groupOpt, minus(allGroupSelectionScopeOpts, groupOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, groupOpt, minus(allConsumerGroupLevelOpts, describeOpt, deleteOpt, resetOffsetsOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, topicOpt, minus(allConsumerGroupLevelOpts, deleteOpt, resetOffsetsOpt)); + } + + @SuppressWarnings("unchecked") + private static Set minus(Set set, T...toRemove) { + Set res = new HashSet<>(set); + for (T t : toRemove) + res.remove(t); + return res; + } +} diff --git a/tools/src/main/java/org/apache/kafka/tools/consumergroup/GroupState.java b/tools/src/main/java/org/apache/kafka/tools/consumergroup/GroupState.java new file mode 100644 index 0000000000000..f90257c87a17b --- /dev/null +++ b/tools/src/main/java/org/apache/kafka/tools/consumergroup/GroupState.java @@ -0,0 +1,35 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.tools.consumergroup; + +import org.apache.kafka.common.Node; + +class GroupState { + public final String group; + public final Node coordinator; + public final String assignmentStrategy; + public final String state; + public final int numMembers; + + public GroupState(String group, Node coordinator, String assignmentStrategy, String state, int numMembers) { + this.group = group; + this.coordinator = coordinator; + this.assignmentStrategy = assignmentStrategy; + this.state = state; + this.numMembers = numMembers; + } +} diff --git a/tools/src/main/java/org/apache/kafka/tools/consumergroup/MemberAssignmentState.java b/tools/src/main/java/org/apache/kafka/tools/consumergroup/MemberAssignmentState.java new file mode 100644 index 0000000000000..77f413d07fafa --- /dev/null +++ b/tools/src/main/java/org/apache/kafka/tools/consumergroup/MemberAssignmentState.java @@ -0,0 +1,42 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.tools.consumergroup; + +import org.apache.kafka.common.TopicPartition; + +import java.util.List; + +class MemberAssignmentState { + public final String group; + public final String consumerId; + public final String host; + public final String clientId; + public final String groupInstanceId; + public final int numPartitions; + public final List assignment; + + public MemberAssignmentState(String group, String consumerId, String host, String clientId, String groupInstanceId, + int numPartitions, List assignment) { + this.group = group; + this.consumerId = consumerId; + this.host = host; + this.clientId = clientId; + this.groupInstanceId = groupInstanceId; + this.numPartitions = numPartitions; + this.assignment = assignment; + } +} diff --git a/tools/src/main/java/org/apache/kafka/tools/consumergroup/PartitionAssignmentState.java b/tools/src/main/java/org/apache/kafka/tools/consumergroup/PartitionAssignmentState.java new file mode 100644 index 0000000000000..27d1f976b9aab --- /dev/null +++ b/tools/src/main/java/org/apache/kafka/tools/consumergroup/PartitionAssignmentState.java @@ -0,0 +1,50 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.tools.consumergroup; + +import org.apache.kafka.common.Node; + +import java.util.Optional; + +class PartitionAssignmentState { + public final String group; + public final Optional coordinator; + public final Optional topic; + public final Optional partition; + public final Optional offset; + public final Optional lag; + public final Optional consumerId; + public final Optional host; + public final Optional clientId; + public final Optional logEndOffset; + + public PartitionAssignmentState(String group, Optional coordinator, Optional topic, + Optional partition, Optional offset, Optional lag, + Optional consumerId, Optional host, Optional clientId, + Optional logEndOffset) { + this.group = group; + this.coordinator = coordinator; + this.topic = topic; + this.partition = partition; + this.offset = offset; + this.lag = lag; + this.consumerId = consumerId; + this.host = host; + this.clientId = clientId; + this.logEndOffset = logEndOffset; + } +} \ No newline at end of file From d62882ddddca1b67493a6322adfc085dded7fd7e Mon Sep 17 00:00:00 2001 From: Nikolay Date: Wed, 6 Dec 2023 22:24:02 +0300 Subject: [PATCH 02/18] Update clients/src/main/java/org/apache/kafka/common/utils/Utils.java Co-authored-by: Taras Ledkov --- clients/src/main/java/org/apache/kafka/common/utils/Utils.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 44dec563d4364..f5500aa2f2b39 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -599,8 +599,7 @@ public static String join(T[] strs, String separator) { * @return The string representation. */ public static String join(Collection collection) { - Objects.requireNonNull(collection); - return mkString(collection.stream(), "", "", ", "); + return join(collection, ", "); } /** From 295d45e4d94c7aaf49764d775fb34f13c10c2c75 Mon Sep 17 00:00:00 2001 From: Nikolay Date: Thu, 7 Dec 2023 12:32:49 +0300 Subject: [PATCH 03/18] Update tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java Co-authored-by: Taras Ledkov --- .../kafka/tools/consumergroup/ConsumerGroupCommandOptions.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java b/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java index 8ab65746e7b06..6ada71a55b789 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java @@ -40,7 +40,7 @@ public class ConsumerGroupCommandOptions extends CommandDefaultOptions { public static final String LIST_DOC = "List all consumer groups."; public static final String DESCRIBE_DOC = "Describe consumer group and list offset lag (number of messages not yet processed) related to given group."; public static final String ALL_GROUPS_DOC = "Apply to all consumer groups."; - public static final String NL = System.getProperty("line.separator"); + public static final String NL = System.lineSeparator(); public static final String DELETE_DOC = "Pass in groups to delete topic partition offsets and ownership information " + "over the entire consumer group. For instance --group g1 --group g2"; public static final String TIMEOUT_MS_DOC = "The timeout that can be set for some use cases. For example, it can be used when describing the group " + From 2e81fa931cfed9138661d5603f81b0d38ba7a466 Mon Sep 17 00:00:00 2001 From: nizhikov Date: Thu, 18 Jan 2024 00:49:14 +0300 Subject: [PATCH 04/18] KAFKA-14589 Code review fixes --- .../org/apache/kafka/common/utils/Utils.java | 14 ++++++++++++ .../ConsumerGroupCommandOptions.java | 22 +++++++------------ .../PartitionAssignmentState.java | 2 +- 3 files changed, 23 insertions(+), 15 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 8ac1ceed8503f..84aeb91df8d74 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -1501,6 +1501,20 @@ public static Set diff(final Supplier> constructor, final Set l return result; } + /** + * @param set Source set. + * @param toRemove Elements to remove. + * @return {@code set} copy without {@code toRemove} elements. + * @param Element type. + */ + @SuppressWarnings("unchecked") + public static Set minus(Set set, T...toRemove) { + Set res = new HashSet<>(set); + for (T t : toRemove) + res.remove(t); + return res; + } + public static Map filterMap(final Map map, final Predicate> filterPredicate) { return map.entrySet().stream().filter(filterPredicate).collect(Collectors.toMap(Entry::getKey, Entry::getValue)); } diff --git a/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java b/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java index 6ada71a55b789..6b0b6056d191d 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java @@ -17,7 +17,6 @@ package org.apache.kafka.tools.consumergroup; import joptsimple.OptionSpec; -import org.apache.kafka.common.utils.Utils; import org.apache.kafka.server.util.CommandDefaultOptions; import org.apache.kafka.server.util.CommandLineUtils; import org.slf4j.Logger; @@ -28,6 +27,9 @@ import java.util.List; import java.util.Set; +import static org.apache.kafka.common.utils.Utils.join; +import static org.apache.kafka.common.utils.Utils.minus; + public class ConsumerGroupCommandOptions extends CommandDefaultOptions { public static final Logger LOGGER = LoggerFactory.getLogger(ConsumerGroupCommandOptions.class); @@ -195,11 +197,11 @@ public void checkArgs() { if (options.has(describeOpt)) { if (!options.has(groupOpt) && !options.has(allGroupsOpt)) CommandLineUtils.printUsageAndExit(parser, - "Option $describeOpt takes one of these options: " + Utils.join(allGroupSelectionScopeOpts)); + "Option " + describeOpt + " takes one of these options: " + join(allGroupSelectionScopeOpts)); List> mutuallyExclusiveOpts = Arrays.asList(membersOpt, offsetsOpt, stateOpt); if (mutuallyExclusiveOpts.stream().mapToInt(o -> options.has(o) ? 1 : 0).sum() > 1) { CommandLineUtils.printUsageAndExit(parser, - "Option " + describeOpt + " takes at most one of these options: " + Utils.join(mutuallyExclusiveOpts)); + "Option " + describeOpt + " takes at most one of these options: " + join(mutuallyExclusiveOpts)); } if (options.has(stateOpt) && options.valueOf(stateOpt) != null) CommandLineUtils.printUsageAndExit(parser, @@ -212,7 +214,7 @@ public void checkArgs() { if (options.has(deleteOpt)) { if (!options.has(groupOpt) && !options.has(allGroupsOpt)) CommandLineUtils.printUsageAndExit(parser, - "Option " + deleteOpt + " takes one of these options: " + Utils.join(allGroupSelectionScopeOpts)); + "Option " + deleteOpt + " takes one of these options: " + join(allGroupSelectionScopeOpts)); if (options.has(topicOpt)) CommandLineUtils.printUsageAndExit(parser, "The consumer does not support topic-specific offset " + "deletion from a consumer group."); @@ -221,7 +223,7 @@ public void checkArgs() { if (options.has(deleteOffsetsOpt)) { if (!options.has(groupOpt) || !options.has(topicOpt)) CommandLineUtils.printUsageAndExit(parser, - "Option " + deleteOffsetsOpt + " takes the following options: " + Utils.join(allDeleteOffsetsOpts)); + "Option " + deleteOffsetsOpt + " takes the following options: " + join(allDeleteOffsetsOpts)); } if (options.has(resetOffsetsOpt)) { @@ -237,7 +239,7 @@ public void checkArgs() { if (!options.has(groupOpt) && !options.has(allGroupsOpt)) CommandLineUtils.printUsageAndExit(parser, - "Option " + resetOffsetsOpt + " takes one of these options: " + Utils.join(allGroupSelectionScopeOpts)); + "Option " + resetOffsetsOpt + " takes one of these options: " + join(allGroupSelectionScopeOpts)); CommandLineUtils.checkInvalidArgs(parser, options, resetToOffsetOpt, minus(allResetOffsetScenarioOpts, resetToOffsetOpt)); CommandLineUtils.checkInvalidArgs(parser, options, resetToDatetimeOpt, minus(allResetOffsetScenarioOpts, resetToDatetimeOpt)); CommandLineUtils.checkInvalidArgs(parser, options, resetByDurationOpt, minus(allResetOffsetScenarioOpts, resetByDurationOpt)); @@ -252,12 +254,4 @@ public void checkArgs() { CommandLineUtils.checkInvalidArgs(parser, options, groupOpt, minus(allConsumerGroupLevelOpts, describeOpt, deleteOpt, resetOffsetsOpt)); CommandLineUtils.checkInvalidArgs(parser, options, topicOpt, minus(allConsumerGroupLevelOpts, deleteOpt, resetOffsetsOpt)); } - - @SuppressWarnings("unchecked") - private static Set minus(Set set, T...toRemove) { - Set res = new HashSet<>(set); - for (T t : toRemove) - res.remove(t); - return res; - } } diff --git a/tools/src/main/java/org/apache/kafka/tools/consumergroup/PartitionAssignmentState.java b/tools/src/main/java/org/apache/kafka/tools/consumergroup/PartitionAssignmentState.java index 27d1f976b9aab..918aaf1cd9508 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumergroup/PartitionAssignmentState.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumergroup/PartitionAssignmentState.java @@ -47,4 +47,4 @@ public PartitionAssignmentState(String group, Optional coordinator, Option this.clientId = clientId; this.logEndOffset = logEndOffset; } -} \ No newline at end of file +} From 70a0a4fc691c4504d94145ee0ea63b7582b20efb Mon Sep 17 00:00:00 2001 From: nizhikov Date: Thu, 18 Jan 2024 01:54:49 +0300 Subject: [PATCH 05/18] KAFKA-14589 Rename partition --- .../group}/ConsumerGroupCommandOptions.java | 2 +- .../tools/{consumergroup => consumer/group}/GroupState.java | 2 +- .../group}/MemberAssignmentState.java | 2 +- .../group}/PartitionAssignmentState.java | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) rename tools/src/main/java/org/apache/kafka/tools/{consumergroup => consumer/group}/ConsumerGroupCommandOptions.java (99%) rename tools/src/main/java/org/apache/kafka/tools/{consumergroup => consumer/group}/GroupState.java (96%) rename tools/src/main/java/org/apache/kafka/tools/{consumergroup => consumer/group}/MemberAssignmentState.java (97%) rename tools/src/main/java/org/apache/kafka/tools/{consumergroup => consumer/group}/PartitionAssignmentState.java (97%) diff --git a/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandOptions.java similarity index 99% rename from tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java rename to tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandOptions.java index 6b0b6056d191d..462b0f25a34ec 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumergroup/ConsumerGroupCommandOptions.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandOptions.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.tools.consumergroup; +package org.apache.kafka.tools.consumer.group; import joptsimple.OptionSpec; import org.apache.kafka.server.util.CommandDefaultOptions; diff --git a/tools/src/main/java/org/apache/kafka/tools/consumergroup/GroupState.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/GroupState.java similarity index 96% rename from tools/src/main/java/org/apache/kafka/tools/consumergroup/GroupState.java rename to tools/src/main/java/org/apache/kafka/tools/consumer/group/GroupState.java index f90257c87a17b..c160f5acadf4f 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumergroup/GroupState.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/GroupState.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.tools.consumergroup; +package org.apache.kafka.tools.consumer.group; import org.apache.kafka.common.Node; diff --git a/tools/src/main/java/org/apache/kafka/tools/consumergroup/MemberAssignmentState.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/MemberAssignmentState.java similarity index 97% rename from tools/src/main/java/org/apache/kafka/tools/consumergroup/MemberAssignmentState.java rename to tools/src/main/java/org/apache/kafka/tools/consumer/group/MemberAssignmentState.java index 77f413d07fafa..040cb1c741ec0 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumergroup/MemberAssignmentState.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/MemberAssignmentState.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.tools.consumergroup; +package org.apache.kafka.tools.consumer.group; import org.apache.kafka.common.TopicPartition; diff --git a/tools/src/main/java/org/apache/kafka/tools/consumergroup/PartitionAssignmentState.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/PartitionAssignmentState.java similarity index 97% rename from tools/src/main/java/org/apache/kafka/tools/consumergroup/PartitionAssignmentState.java rename to tools/src/main/java/org/apache/kafka/tools/consumer/group/PartitionAssignmentState.java index 918aaf1cd9508..eec449ed557ed 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumergroup/PartitionAssignmentState.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/PartitionAssignmentState.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.tools.consumergroup; +package org.apache.kafka.tools.consumer.group; import org.apache.kafka.common.Node; From 54986b54255e0e0e316fab3f6d7bba0cf03de316 Mon Sep 17 00:00:00 2001 From: nizhikov Date: Thu, 18 Jan 2024 19:22:38 +0300 Subject: [PATCH 06/18] KAFKA-14589 Review fixes --- .../org/apache/kafka/common/utils/Utils.java | 23 ------------------- .../org/apache/kafka/tools/ToolsUtils.java | 15 ++++++++++++ .../group/ConsumerGroupCommandOptions.java | 12 +++++----- 3 files changed, 21 insertions(+), 29 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 84aeb91df8d74..c316b7a181662 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -594,15 +594,6 @@ public static String join(T[] strs, String separator) { return join(Arrays.asList(strs), separator); } - /** - * Create a string representation of a collection joined by ", ". - * @param collection The list of items - * @return The string representation. - */ - public static String join(Collection collection) { - return join(collection, ", "); - } - /** * Create a string representation of a collection joined by the given separator * @param collection The list of items @@ -1501,20 +1492,6 @@ public static Set diff(final Supplier> constructor, final Set l return result; } - /** - * @param set Source set. - * @param toRemove Elements to remove. - * @return {@code set} copy without {@code toRemove} elements. - * @param Element type. - */ - @SuppressWarnings("unchecked") - public static Set minus(Set set, T...toRemove) { - Set res = new HashSet<>(set); - for (T t : toRemove) - res.remove(t); - return res; - } - public static Map filterMap(final Map map, final Predicate> filterPredicate) { return map.entrySet().stream().filter(filterPredicate).collect(Collectors.toMap(Entry::getKey, Entry::getValue)); } diff --git a/tools/src/main/java/org/apache/kafka/tools/ToolsUtils.java b/tools/src/main/java/org/apache/kafka/tools/ToolsUtils.java index a3eb4ab4bcb6b..9c6a7a0d1ccb6 100644 --- a/tools/src/main/java/org/apache/kafka/tools/ToolsUtils.java +++ b/tools/src/main/java/org/apache/kafka/tools/ToolsUtils.java @@ -140,4 +140,19 @@ public static Set duplicates(List s) { }); return duplicates; } + + /** + * @param set Source set. + * @param toRemove Elements to remove. + * @return {@code set} copy without {@code toRemove} elements. + * @param Element type. + */ + @SuppressWarnings("unchecked") + public static Set minus(Set set, T...toRemove) { + Set res = new HashSet<>(set); + for (T t : toRemove) + res.remove(t); + return res; + } + } diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandOptions.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandOptions.java index 462b0f25a34ec..0448cc4f12e63 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandOptions.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandOptions.java @@ -28,7 +28,7 @@ import java.util.Set; import static org.apache.kafka.common.utils.Utils.join; -import static org.apache.kafka.common.utils.Utils.minus; +import static org.apache.kafka.tools.ToolsUtils.minus; public class ConsumerGroupCommandOptions extends CommandDefaultOptions { public static final Logger LOGGER = LoggerFactory.getLogger(ConsumerGroupCommandOptions.class); @@ -197,11 +197,11 @@ public void checkArgs() { if (options.has(describeOpt)) { if (!options.has(groupOpt) && !options.has(allGroupsOpt)) CommandLineUtils.printUsageAndExit(parser, - "Option " + describeOpt + " takes one of these options: " + join(allGroupSelectionScopeOpts)); + "Option " + describeOpt + " takes one of these options: " + join(allGroupSelectionScopeOpts, ", ")); List> mutuallyExclusiveOpts = Arrays.asList(membersOpt, offsetsOpt, stateOpt); if (mutuallyExclusiveOpts.stream().mapToInt(o -> options.has(o) ? 1 : 0).sum() > 1) { CommandLineUtils.printUsageAndExit(parser, - "Option " + describeOpt + " takes at most one of these options: " + join(mutuallyExclusiveOpts)); + "Option " + describeOpt + " takes at most one of these options: " + join(mutuallyExclusiveOpts, ", ")); } if (options.has(stateOpt) && options.valueOf(stateOpt) != null) CommandLineUtils.printUsageAndExit(parser, @@ -214,7 +214,7 @@ public void checkArgs() { if (options.has(deleteOpt)) { if (!options.has(groupOpt) && !options.has(allGroupsOpt)) CommandLineUtils.printUsageAndExit(parser, - "Option " + deleteOpt + " takes one of these options: " + join(allGroupSelectionScopeOpts)); + "Option " + deleteOpt + " takes one of these options: " + join(allGroupSelectionScopeOpts, ", ")); if (options.has(topicOpt)) CommandLineUtils.printUsageAndExit(parser, "The consumer does not support topic-specific offset " + "deletion from a consumer group."); @@ -223,7 +223,7 @@ public void checkArgs() { if (options.has(deleteOffsetsOpt)) { if (!options.has(groupOpt) || !options.has(topicOpt)) CommandLineUtils.printUsageAndExit(parser, - "Option " + deleteOffsetsOpt + " takes the following options: " + join(allDeleteOffsetsOpts)); + "Option " + deleteOffsetsOpt + " takes the following options: " + join(allDeleteOffsetsOpts, ", ")); } if (options.has(resetOffsetsOpt)) { @@ -239,7 +239,7 @@ public void checkArgs() { if (!options.has(groupOpt) && !options.has(allGroupsOpt)) CommandLineUtils.printUsageAndExit(parser, - "Option " + resetOffsetsOpt + " takes one of these options: " + join(allGroupSelectionScopeOpts)); + "Option " + resetOffsetsOpt + " takes one of these options: " + join(allGroupSelectionScopeOpts, ", ")); CommandLineUtils.checkInvalidArgs(parser, options, resetToOffsetOpt, minus(allResetOffsetScenarioOpts, resetToOffsetOpt)); CommandLineUtils.checkInvalidArgs(parser, options, resetToDatetimeOpt, minus(allResetOffsetScenarioOpts, resetToDatetimeOpt)); CommandLineUtils.checkInvalidArgs(parser, options, resetByDurationOpt, minus(allResetOffsetScenarioOpts, resetByDurationOpt)); From 40e8f6d44252e89de379206053742370713b585d Mon Sep 17 00:00:00 2001 From: nizhikov Date: Thu, 18 Jan 2024 19:27:47 +0300 Subject: [PATCH 07/18] KAFKA-14589 Review fixes --- checkstyle/import-control.xml | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/checkstyle/import-control.xml b/checkstyle/import-control.xml index 3316b593b167e..c39075677b22c 100644 --- a/checkstyle/import-control.xml +++ b/checkstyle/import-control.xml @@ -327,6 +327,10 @@ + + + + From c957587580e36237a756199fcf9fa27d6226971d Mon Sep 17 00:00:00 2001 From: "n.izhikov" Date: Mon, 22 Jan 2024 18:45:26 +0300 Subject: [PATCH 08/18] KAFKA-14589 Review fixes --- .../consumer/group/PartitionAssignmentState.java | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/PartitionAssignmentState.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/PartitionAssignmentState.java index eec449ed557ed..396032f0a0c29 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/PartitionAssignmentState.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/PartitionAssignmentState.java @@ -19,23 +19,25 @@ import org.apache.kafka.common.Node; import java.util.Optional; +import java.util.OptionalInt; +import java.util.OptionalLong; class PartitionAssignmentState { public final String group; public final Optional coordinator; public final Optional topic; - public final Optional partition; - public final Optional offset; - public final Optional lag; + public final OptionalInt partition; + public final OptionalLong offset; + public final OptionalLong lag; public final Optional consumerId; public final Optional host; public final Optional clientId; - public final Optional logEndOffset; + public final OptionalLong logEndOffset; public PartitionAssignmentState(String group, Optional coordinator, Optional topic, - Optional partition, Optional offset, Optional lag, + OptionalInt partition, OptionalLong offset, OptionalLong lag, Optional consumerId, Optional host, Optional clientId, - Optional logEndOffset) { + OptionalLong logEndOffset) { this.group = group; this.coordinator = coordinator; this.topic = topic; From 82058f49783ba4734ddb45d7be41efdeec2e2a60 Mon Sep 17 00:00:00 2001 From: "n.izhikov" Date: Mon, 22 Jan 2024 21:42:49 +0300 Subject: [PATCH 09/18] KAFKA-14589 Review fixes --- .../consumer/group/ConsumerGroupCommandOptions.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandOptions.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandOptions.java index 0448cc4f12e63..045d296444dfb 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandOptions.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandOptions.java @@ -240,14 +240,14 @@ public void checkArgs() { if (!options.has(groupOpt) && !options.has(allGroupsOpt)) CommandLineUtils.printUsageAndExit(parser, "Option " + resetOffsetsOpt + " takes one of these options: " + join(allGroupSelectionScopeOpts, ", ")); - CommandLineUtils.checkInvalidArgs(parser, options, resetToOffsetOpt, minus(allResetOffsetScenarioOpts, resetToOffsetOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, resetToOffsetOpt, minus(allResetOffsetScenarioOpts, resetToOffsetOpt)); CommandLineUtils.checkInvalidArgs(parser, options, resetToDatetimeOpt, minus(allResetOffsetScenarioOpts, resetToDatetimeOpt)); CommandLineUtils.checkInvalidArgs(parser, options, resetByDurationOpt, minus(allResetOffsetScenarioOpts, resetByDurationOpt)); CommandLineUtils.checkInvalidArgs(parser, options, resetToEarliestOpt, minus(allResetOffsetScenarioOpts, resetToEarliestOpt)); - CommandLineUtils.checkInvalidArgs(parser, options, resetToLatestOpt, minus(allResetOffsetScenarioOpts, resetToLatestOpt)); - CommandLineUtils.checkInvalidArgs(parser, options, resetToCurrentOpt, minus(allResetOffsetScenarioOpts, resetToCurrentOpt)); - CommandLineUtils.checkInvalidArgs(parser, options, resetShiftByOpt, minus(allResetOffsetScenarioOpts, resetShiftByOpt)); - CommandLineUtils.checkInvalidArgs(parser, options, resetFromFileOpt, minus(allResetOffsetScenarioOpts, resetFromFileOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, resetToLatestOpt, minus(allResetOffsetScenarioOpts, resetToLatestOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, resetToCurrentOpt, minus(allResetOffsetScenarioOpts, resetToCurrentOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, resetShiftByOpt, minus(allResetOffsetScenarioOpts, resetShiftByOpt)); + CommandLineUtils.checkInvalidArgs(parser, options, resetFromFileOpt, minus(allResetOffsetScenarioOpts, resetFromFileOpt)); } CommandLineUtils.checkInvalidArgs(parser, options, groupOpt, minus(allGroupSelectionScopeOpts, groupOpt)); From aed456764e890eb3f61c6671faa2e8a41a6b5753 Mon Sep 17 00:00:00 2001 From: "n.izhikov" Date: Tue, 23 Jan 2024 12:58:31 +0300 Subject: [PATCH 10/18] KAFKA-14589 WIP --- .../kafka/admin/ConsumerGroupCommand.scala | 2 +- .../kafka/tools/{reassign => }/Tuple2.java | 2 +- .../reassign/ReassignPartitionsCommand.java | 1 + .../group/ConsumerGroupCommandTest.java | 312 ++++++++++++++++++ .../group/ConsumerGroupServiceTest.java | 287 ++++++++++++++++ .../ReassignPartitionsIntegrationTest.java | 1 + .../reassign/ReassignPartitionsUnitTest.java | 1 + 7 files changed, 604 insertions(+), 2 deletions(-) rename tools/src/main/java/org/apache/kafka/tools/{reassign => }/Tuple2.java (97%) create mode 100644 tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandTest.java create mode 100644 tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java diff --git a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala index 7db2707e287e3..600e053e5001c 100755 --- a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala +++ b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala @@ -125,7 +125,7 @@ object ConsumerGroupCommand extends Logging { } } - private[admin] case class PartitionAssignmentState(group: String, coordinator: Option[Node], topic: Option[String], + case class PartitionAssignmentState(group: String, coordinator: Option[Node], topic: Option[String], partition: Option[Int], offset: Option[Long], lag: Option[Long], consumerId: Option[String], host: Option[String], clientId: Option[String], logEndOffset: Option[Long]) diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/Tuple2.java b/tools/src/main/java/org/apache/kafka/tools/Tuple2.java similarity index 97% rename from tools/src/main/java/org/apache/kafka/tools/reassign/Tuple2.java rename to tools/src/main/java/org/apache/kafka/tools/Tuple2.java index a9e84317cf718..02dd4bf3981a8 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/Tuple2.java +++ b/tools/src/main/java/org/apache/kafka/tools/Tuple2.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.tools.reassign; +package org.apache.kafka.tools; import java.util.Objects; diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java b/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java index 6a2d9b9346863..3623346980743 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java @@ -49,6 +49,7 @@ import org.apache.kafka.server.util.json.JsonValue; import org.apache.kafka.tools.TerseException; import org.apache.kafka.tools.ToolsUtils; +import org.apache.kafka.tools.Tuple2; import java.io.IOException; import java.util.ArrayList; diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandTest.java new file mode 100644 index 0000000000000..6c3712a5d218e --- /dev/null +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandTest.java @@ -0,0 +1,312 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.tools.consumer.group; + +import kafka.admin.ConsumerGroupCommand; +import kafka.server.KafkaConfig; +import kafka.utils.TestUtils; +import org.apache.kafka.clients.admin.AdminClientConfig; +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.clients.consumer.RangeAssignor; +import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.WakeupException; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.TestInfo; +import scala.collection.JavaConverters; +import scala.collection.Seq; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.Properties; +import java.util.Set; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; +import java.util.stream.IntStream; + +public class ConsumerGroupCommandTest extends kafka.integration.KafkaServerTestHarness { + public static final String TOPIC = "foo"; + public static final String GROUP = "test.group"; + + List consumerGroupService = new ArrayList<>(); + List consumerGroupExecutors = new ArrayList<>(); + + @Override + public Seq generateConfigs() { + List cfgs = new ArrayList<>(); + + TestUtils.createBrokerConfigs( + 1, + zkConnectOrNull(), + false, + true, + scala.None$.empty(), + scala.None$.empty(), + scala.None$.empty(), + true, + false, + false, + false, + scala.collection.immutable.Map$.MODULE$.empty(), + 1, + false, + 1, + (short) 1, + 0, + false + ).foreach(props -> { + cfgs.add(KafkaConfig.fromProps(props)); + return null; + }); + + return seq(cfgs); + } + + @BeforeEach + @Override + public void setUp(TestInfo testInfo) { + super.setUp(testInfo); + createTopic(TOPIC, 1, 1, new Properties(), listenerName(), new Properties()); + } + + @AfterEach + @Override + public void tearDown() { + consumerGroupService.forEach(ConsumerGroupCommand.ConsumerGroupService::close); + consumerGroupExecutors.forEach(AbstractConsumerGroupExecutor::shutdown); + super.tearDown(); + } + + Map committedOffsets(String topic, String group) { + try (Consumer consumer = createNoAutoCommitConsumer(group)) { + Set partitions = consumer.partitionsFor(topic).stream() + .map(partitionInfo -> new TopicPartition(partitionInfo.topic(), partitionInfo.partition())) + .collect(Collectors.toSet()); + return consumer.committed(partitions).entrySet().stream() + .filter(e -> e.getValue() != null) + .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().offset())); + } + } + + Consumer createNoAutoCommitConsumer(String group) { + Properties props = new Properties(); + props.put("bootstrap.servers", bootstrapServers(listenerName())); + props.put("group.id", group); + props.put("enable.auto.commit", "false"); + return new KafkaConsumer<>(props, new StringDeserializer(), new StringDeserializer()); + } + + ConsumerGroupCommand.ConsumerGroupService getConsumerGroupService(String[] args) { + ConsumerGroupCommand.ConsumerGroupCommandOptions opts = new ConsumerGroupCommand.ConsumerGroupCommandOptions(args); + ConsumerGroupCommand.ConsumerGroupService service = new ConsumerGroupCommand.ConsumerGroupService( + opts, + JavaConverters.asScala(Collections.singletonMap(AdminClientConfig.RETRIES_CONFIG, Integer.toString(Integer.MAX_VALUE))) + ); + + consumerGroupService.add(0, service); + return service; + } + + ConsumerGroupExecutor addConsumerGroupExecutor(int numConsumers) { + return addConsumerGroupExecutor(numConsumers, TOPIC, GROUP, RangeAssignor.class.getName(), Optional.empty(), false); + } + + ConsumerGroupExecutor addConsumerGroupExecutor(int numConsumers, String topic, String group) { + return addConsumerGroupExecutor(numConsumers, topic, group, RangeAssignor.class.getName(), Optional.empty(), false); + } + + ConsumerGroupExecutor addConsumerGroupExecutor(int numConsumers, String topic, String group, String strategy, + Optional customPropsOpt, boolean syncCommit) { + ConsumerGroupExecutor executor = new ConsumerGroupExecutor(bootstrapServers(listenerName()), numConsumers, group, + topic, strategy, customPropsOpt, syncCommit); + addExecutor(executor); + return executor; + } + + SimpleConsumerGroupExecutor addSimpleGroupExecutor(String group) { + return addSimpleGroupExecutor(Arrays.asList(new TopicPartition(TOPIC, 0)), group); + } + + SimpleConsumerGroupExecutor addSimpleGroupExecutor(Collection partitions, String group) { + SimpleConsumerGroupExecutor executor = new SimpleConsumerGroupExecutor(bootstrapServers(listenerName()), group, partitions); + addExecutor(executor); + return executor; + } + + private AbstractConsumerGroupExecutor addExecutor(AbstractConsumerGroupExecutor executor) { + consumerGroupExecutors.add(0, executor); + return executor; + } + + abstract class AbstractConsumerRunnable implements Runnable { + final String broker; + final String groupId; + final Optional customPropsOpt; + final boolean syncCommit; + + final Properties props = new Properties(); + KafkaConsumer consumer; + + boolean configured = false; + + public AbstractConsumerRunnable(String broker, String groupId, Optional customPropsOpt, boolean syncCommit) { + this.broker = broker; + this.groupId = groupId; + this.customPropsOpt = customPropsOpt; + this.syncCommit = syncCommit; + } + + void configure() { + configured = true; + configure(props); + customPropsOpt.ifPresent(props::putAll); + consumer = new KafkaConsumer<>(props); + } + + void configure(Properties props) { + props.put("bootstrap.servers", broker); + props.put("group.id", groupId); + props.put("key.deserializer", StringDeserializer.class.getName()); + props.put("value.deserializer", StringDeserializer.class.getName()); + } + + abstract void subscribe(); + + @Override + public void run() { + assert configured : "Must call configure before use"; + try { + subscribe(); + while (true) { + consumer.poll(Duration.ofMillis(Long.MAX_VALUE)); + if (syncCommit) + consumer.commitSync(); + } + } catch (WakeupException e) { + // OK + } finally { + consumer.close(); + } + } + + void shutdown() { + consumer.wakeup(); + } + } + + class ConsumerRunnable extends AbstractConsumerRunnable { + final String topic; + final String strategy; + + public ConsumerRunnable(String broker, String groupId, String topic, String strategy, + Optional customPropsOpt, boolean syncCommit) { + super(broker, groupId, customPropsOpt, syncCommit); + + this.topic = topic; + this.strategy = strategy; + } + + @Override + void configure(Properties props) { + super.configure(props); + props.put("partition.assignment.strategy", strategy); + } + + @Override + void subscribe() { + consumer.subscribe(Collections.singleton(topic)); + } + } + + class SimpleConsumerRunnable extends AbstractConsumerRunnable { + final Collection partitions; + + public SimpleConsumerRunnable(String broker, String groupId, Collection partitions) { + super(broker, groupId, Optional.empty(), false); + + this.partitions = partitions; + } + + @Override + void subscribe() { + consumer.assign(partitions); + } + } + + class AbstractConsumerGroupExecutor { + final int numThreads; + final ExecutorService executor; + final List consumers = new ArrayList<>(); + + public AbstractConsumerGroupExecutor(int numThreads) { + this.numThreads = numThreads; + this.executor = Executors.newFixedThreadPool(numThreads); + } + + void submit(AbstractConsumerRunnable consumerThread) { + consumers.add(consumerThread); + executor.submit(consumerThread); + } + + void shutdown() { + consumers.forEach(AbstractConsumerRunnable::shutdown); + executor.shutdown(); + try { + executor.awaitTermination(5000, TimeUnit.MILLISECONDS); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } + } + } + + class ConsumerGroupExecutor extends AbstractConsumerGroupExecutor { + public ConsumerGroupExecutor(String broker, int numConsumers, String groupId, String topic, String strategy, + Optional customPropsOpt, boolean syncCommit) { + super(numConsumers); + IntStream.rangeClosed(1, numConsumers).forEach(i -> { + ConsumerRunnable th = new ConsumerRunnable(broker, groupId, topic, strategy, customPropsOpt, syncCommit); + th.configure(); + submit(th); + }); + } + } + + class SimpleConsumerGroupExecutor extends AbstractConsumerGroupExecutor { + public SimpleConsumerGroupExecutor(String broker, String groupId, Collection partitions) { + super(1); + + SimpleConsumerRunnable th = new SimpleConsumerRunnable(broker, groupId, partitions); + th.configure(); + submit(th); + } + } + + @SuppressWarnings({"deprecation"}) + static Seq seq(Collection seq) { + return JavaConverters.asScalaIteratorConverter(seq.iterator()).asScala().toSeq(); + } +} \ No newline at end of file diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java new file mode 100644 index 0000000000000..7cc46acff118c --- /dev/null +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java @@ -0,0 +1,287 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.tools.consumer.group; + +import kafka.admin.ConsumerGroupCommand; +import org.apache.kafka.clients.admin.Admin; +import org.apache.kafka.clients.admin.AdminClientTestUtils; +import org.apache.kafka.clients.admin.ConsumerGroupDescription; +import org.apache.kafka.clients.admin.DescribeConsumerGroupsResult; +import org.apache.kafka.clients.admin.DescribeTopicsResult; +import org.apache.kafka.clients.admin.ListConsumerGroupOffsetsResult; +import org.apache.kafka.clients.admin.ListConsumerGroupOffsetsSpec; +import org.apache.kafka.clients.admin.ListOffsetsResult; +import org.apache.kafka.clients.admin.ListOffsetsResult.ListOffsetsResultInfo; +import org.apache.kafka.clients.admin.MemberAssignment; +import org.apache.kafka.clients.admin.MemberDescription; +import org.apache.kafka.clients.admin.OffsetSpec; +import org.apache.kafka.clients.admin.TopicDescription; +import org.apache.kafka.clients.consumer.OffsetAndMetadata; +import org.apache.kafka.clients.consumer.RangeAssignor; +import org.apache.kafka.common.ConsumerGroupState; +import org.apache.kafka.common.KafkaFuture; +import org.apache.kafka.common.Node; +import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.TopicPartitionInfo; +import org.apache.kafka.common.internals.KafkaFutureImpl; +import org.apache.kafka.common.utils.Utils; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentMatcher; +import org.mockito.ArgumentMatchers; +import scala.Option; +import scala.Tuple2; +import scala.collection.JavaConverters; +import scala.collection.Seq; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.Objects; +import java.util.Set; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.function.Function; +import java.util.stream.Collectors; +import java.util.stream.IntStream; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +public class ConsumerGroupServiceTest { + public static final String GROUP = "testGroup"; + + public static final int NUM_PARTITIONS = 10; + + private static final List TOPICS = IntStream.range(0, 5).mapToObj(i -> "testTopic" + i).collect(Collectors.toList()); + + private static final List TOPIC_PARTITIONS = TOPICS.stream() + .flatMap(topic -> IntStream.range(0, NUM_PARTITIONS).mapToObj(i -> new TopicPartition(topic, i))) + .collect(Collectors.toList()); + + private Admin admin = mock(Admin.class); + + @Test + public void testAdminRequestsForDescribeOffsets() throws Exception { + String[] args = new String[]{"--bootstrap-server", "localhost:9092", "--group", GROUP, "--describe", "--offsets"}; + ConsumerGroupCommand.ConsumerGroupService groupService = consumerGroupService(args); + + when(admin.describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(GROUP)), any())) + .thenReturn(describeGroupsResult(ConsumerGroupState.STABLE)); + when(admin.listConsumerGroupOffsets(ArgumentMatchers.eq(listConsumerGroupOffsetsSpec()), any())) + .thenReturn(listGroupOffsetsResult(GROUP)); + when(admin.listOffsets(offsetsArgMatcher(), any())) + .thenReturn(listOffsetsResult()); + + Tuple2, Option>> res = groupService.collectGroupOffsets(GROUP); + assertEquals(Optional.of("Stable"), res._1); + assertTrue(res._2.isDefined()); + assertEquals(TOPIC_PARTITIONS.size(), res._2.get().size()); + + verify(admin, times(1)).describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(GROUP)), any()); + verify(admin, times(1)).listConsumerGroupOffsets(ArgumentMatchers.eq(listConsumerGroupOffsetsSpec()), any()); + verify(admin, times(1)).listOffsets(offsetsArgMatcher(), any()); + } + + @Test + public void testAdminRequestsForDescribeNegativeOffsets() throws Exception { + String[] args = new String[]{"--bootstrap-server", "localhost:9092", "--group", GROUP, "--describe", "--offsets"}; + ConsumerGroupCommand.ConsumerGroupService groupService = consumerGroupService(args); + + TopicPartition testTopicPartition0 = new TopicPartition("testTopic1", 0); + TopicPartition testTopicPartition1 = new TopicPartition("testTopic1", 1); + TopicPartition testTopicPartition2 = new TopicPartition("testTopic1", 2); + TopicPartition testTopicPartition3 = new TopicPartition("testTopic2", 0); + TopicPartition testTopicPartition4 = new TopicPartition("testTopic2", 1); + TopicPartition testTopicPartition5 = new TopicPartition("testTopic2", 2); + + // Some topic's partitions gets valid OffsetAndMetadata values, other gets nulls values (negative integers) and others aren't defined + Map committedOffsets = new HashMap<>(); + + committedOffsets.put(testTopicPartition1, new OffsetAndMetadata(100)); + committedOffsets.put(testTopicPartition2, null); + committedOffsets.put(testTopicPartition3, new OffsetAndMetadata(100)); + committedOffsets.put(testTopicPartition4, new OffsetAndMetadata(100)); + committedOffsets.put(testTopicPartition5, null); + + ListOffsetsResultInfo resultInfo = new ListOffsetsResultInfo(100, System.currentTimeMillis(), Optional.of(1)); + Map> endOffsets = new HashMap<>(); + + endOffsets.put(testTopicPartition0, KafkaFuture.completedFuture(resultInfo)); + endOffsets.put(testTopicPartition1, KafkaFuture.completedFuture(resultInfo)); + endOffsets.put(testTopicPartition2, KafkaFuture.completedFuture(resultInfo)); + endOffsets.put(testTopicPartition3, KafkaFuture.completedFuture(resultInfo)); + endOffsets.put(testTopicPartition4, KafkaFuture.completedFuture(resultInfo)); + endOffsets.put(testTopicPartition5, KafkaFuture.completedFuture(resultInfo)); + + Set assignedTopicPartitions = new HashSet<>(Arrays.asList(testTopicPartition0, testTopicPartition1, testTopicPartition2)); + Set unassignedTopicPartitions = new HashSet<>(Arrays.asList(testTopicPartition3, testTopicPartition4, testTopicPartition5)); + + ConsumerGroupDescription consumerGroupDescription = new ConsumerGroupDescription(GROUP, + true, + Collections.singleton(new MemberDescription("member1", Optional.of("instance1"), "client1", "host1", new MemberAssignment(assignedTopicPartitions))), + RangeAssignor.class.getName(), + ConsumerGroupState.STABLE, + new Node(1, "localhost", 9092)); + + Function, ArgumentMatcher>> offsetsArgMatcher = expectedPartitions -> + topicPartitionOffsets -> topicPartitionOffsets != null && topicPartitionOffsets.keySet().equals(expectedPartitions); + + KafkaFutureImpl future = new KafkaFutureImpl<>(); + future.complete(consumerGroupDescription); + when(admin.describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(GROUP)), any())) + .thenReturn(new DescribeConsumerGroupsResult(Collections.singletonMap(GROUP, future))); + when(admin.listConsumerGroupOffsets(ArgumentMatchers.eq(listConsumerGroupOffsetsSpec()), any())) + .thenReturn( + AdminClientTestUtils.listConsumerGroupOffsetsResult( + Collections.singletonMap(GROUP, committedOffsets))); + when(admin.listOffsets( + ArgumentMatchers.argThat(offsetsArgMatcher.apply(assignedTopicPartitions)), + any() + )).thenReturn(new ListOffsetsResult(endOffsets.entrySet().stream().filter(e -> assignedTopicPartitions.contains(e.getKey())) + .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)))); + when(admin.listOffsets( + ArgumentMatchers.argThat(offsetsArgMatcher.apply(unassignedTopicPartitions)), + any() + )).thenReturn(new ListOffsetsResult(endOffsets.entrySet().stream().filter(e -> unassignedTopicPartitions.contains(e.getKey())) + .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)))); + + Tuple2, Option>> res = groupService.collectGroupOffsets(GROUP); + Option state = res._1; + Option> assignments = res._2; + + Map> returnedOffsets = assignments.map(results -> + results.stream().collect(Collectors.toMap( + assignment -> new TopicPartition(assignment.topic.get(), assignment.partition.get()), + assignment -> assignment.offset)) + ).orElse(Collections.emptyMap()); + + Map> expectedOffsets = new HashMap<>(); + + expectedOffsets.put(testTopicPartition0, Optional.empty()); + expectedOffsets.put(testTopicPartition1, Optional.of(100L)); + expectedOffsets.put(testTopicPartition2, Optional.empty()); + expectedOffsets.put(testTopicPartition3, Optional.of(100L)); + expectedOffsets.put(testTopicPartition4, Optional.of(100L)); + expectedOffsets.put(testTopicPartition5, Optional.empty()); + + assertEquals(Optional.of("Stable"), state); + assertEquals(expectedOffsets, returnedOffsets); + + verify(admin, times(1)).describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(GROUP)), any()); + verify(admin, times(1)).listConsumerGroupOffsets(ArgumentMatchers.eq(listConsumerGroupOffsetsSpec()), any()); + verify(admin, times(1)).listOffsets(ArgumentMatchers.argThat(offsetsArgMatcher.apply(assignedTopicPartitions)), any()); + verify(admin, times(1)).listOffsets(ArgumentMatchers.argThat(offsetsArgMatcher.apply(unassignedTopicPartitions)), any()); + } + + @Test + public void testAdminRequestsForResetOffsets() { + List args = new ArrayList<>(Arrays.asList("--bootstrap-server", "localhost:9092", "--group", GROUP, "--reset-offsets", "--to-latest")); + List topicsWithoutPartitionsSpecified = TOPICS.subList(1, TOPICS.size()); + List topicArgs = new ArrayList<>(Arrays.asList("--topic", TOPICS.get(0) + ":" + Utils.mkString(IntStream.range(0, NUM_PARTITIONS).mapToObj(Integer::toString), "", "", ","))); + topicsWithoutPartitionsSpecified.forEach(topic -> topicArgs.addAll(Arrays.asList("--topic", topic))); + + args.addAll(topicArgs); + ConsumerGroupCommand.ConsumerGroupService groupService = consumerGroupService(args.toArray(new String[0])); + + when(admin.describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(GROUP)), any())) + .thenReturn(describeGroupsResult(ConsumerGroupState.DEAD)); + when(admin.describeTopics(ArgumentMatchers.eq(topicsWithoutPartitionsSpecified), any())) + .thenReturn(describeTopicsResult(topicsWithoutPartitionsSpecified)); + when(admin.listOffsets(offsetsArgMatcher(), any())) + .thenReturn(listOffsetsResult()); + + scala.collection.Map> resetResult = groupService.resetOffsets(); + assertEquals(Collections.singleton(GROUP), resetResult.keySet()); + assertEquals(new HashSet<>(TOPIC_PARTITIONS), resetResult.get(GROUP).get().keys()); + + verify(admin, times(1)).describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(GROUP)), any()); + verify(admin, times(1)).describeTopics(ArgumentMatchers.eq(topicsWithoutPartitionsSpecified), any()); + verify(admin, times(1)).listOffsets(offsetsArgMatcher(), any()); + } + + private ConsumerGroupCommand.ConsumerGroupService consumerGroupService(String[] args) { + return new ConsumerGroupCommand.ConsumerGroupService(new ConsumerGroupCommandOptions(args), JavaConverters.asScala(Collections.emptyMap())) { + @Override + public Admin createAdminClient(scala.collection.Map configOverrides) { + return admin; + } + }; + } + + private DescribeConsumerGroupsResult describeGroupsResult(ConsumerGroupState groupState) { + MemberDescription member1 = new MemberDescription("member1", Optional.of("instance1"), "client1", "host1", null); + ConsumerGroupDescription description = new ConsumerGroupDescription(GROUP, + true, + Collections.singleton(member1), + RangeAssignor.class.getName(), + groupState, + new Node(1, "localhost", 9092)); + KafkaFutureImpl future = new KafkaFutureImpl<>(); + future.complete(description); + return new DescribeConsumerGroupsResult(Collections.singletonMap(GROUP, future)); + } + + private ListConsumerGroupOffsetsResult listGroupOffsetsResult(String groupId) { + Map offsets = TOPIC_PARTITIONS.stream().collect(Collectors.toMap( + Function.identity(), + tp -> new OffsetAndMetadata(100))); + return AdminClientTestUtils.listConsumerGroupOffsetsResult(Collections.singletonMap(groupId, offsets)); + } + + private Map offsetsArgMatcher() { + Map expectedOffsets = TOPIC_PARTITIONS.stream().collect(Collectors.toMap( + Function.identity(), + tp -> OffsetSpec.latest() + )); + return ArgumentMatchers.argThat(map -> + Objects.equals(map.keySet(), expectedOffsets.keySet()) && map.values().stream().allMatch(v -> v instanceof OffsetSpec.LatestSpec) + ); + } + + private ListOffsetsResult listOffsetsResult() { + ListOffsetsResultInfo resultInfo = new ListOffsetsResultInfo(100, System.currentTimeMillis(), Optional.of(1)); + Map> futures = TOPIC_PARTITIONS.stream().collect(Collectors.toMap( + Function.identity(), + tp -> KafkaFuture.completedFuture(resultInfo))); + return new ListOffsetsResult(futures); + } + + private DescribeTopicsResult describeTopicsResult(Collection topics) { + Map topicDescriptions = new HashMap<>(); + + topics.forEach(topic -> { + List partitions = IntStream.range(0, NUM_PARTITIONS) + .mapToObj(i -> new TopicPartitionInfo(i, null, Collections.emptyList(), Collections.emptyList())) + .collect(Collectors.toList()); + topicDescriptions.put(topic, new TopicDescription(topic, false, partitions)); + }); + return AdminClientTestUtils.describeTopicsResult(topicDescriptions); + } + + private Map listConsumerGroupOffsetsSpec() { + return Collections.singletonMap(GROUP, new ListConsumerGroupOffsetsSpec()); + } +} diff --git a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsIntegrationTest.java b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsIntegrationTest.java index b71331a6b233f..83fc665e3e487 100644 --- a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsIntegrationTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsIntegrationTest.java @@ -43,6 +43,7 @@ import org.apache.kafka.common.utils.Time; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.tools.TerseException; +import org.apache.kafka.tools.Tuple2; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Timeout; import org.junit.jupiter.params.ParameterizedTest; diff --git a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsUnitTest.java b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsUnitTest.java index d0e1cb4baaf69..699e048fb206e 100644 --- a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsUnitTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsUnitTest.java @@ -31,6 +31,7 @@ import org.apache.kafka.common.utils.Time; import org.apache.kafka.server.common.AdminCommandFailedException; import org.apache.kafka.server.common.AdminOperationException; +import org.apache.kafka.tools.Tuple2; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; From 4700e381786c245fe43dd99be7e0f37f2bf8cef5 Mon Sep 17 00:00:00 2001 From: "n.izhikov" Date: Tue, 23 Jan 2024 17:59:41 +0300 Subject: [PATCH 11/18] KAFKA-14589 WIP --- .../admin/ConsumerGroupServiceTest.scala | 226 ------------- .../group/ConsumerGroupCommandTest.java | 312 ------------------ 2 files changed, 538 deletions(-) delete mode 100644 core/src/test/scala/unit/kafka/admin/ConsumerGroupServiceTest.scala delete mode 100644 tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandTest.java diff --git a/core/src/test/scala/unit/kafka/admin/ConsumerGroupServiceTest.scala b/core/src/test/scala/unit/kafka/admin/ConsumerGroupServiceTest.scala deleted file mode 100644 index 4aa30fbeac6ce..0000000000000 --- a/core/src/test/scala/unit/kafka/admin/ConsumerGroupServiceTest.scala +++ /dev/null @@ -1,226 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package kafka.admin - -import java.util -import java.util.{Collections, Optional} - -import kafka.admin.ConsumerGroupCommand.{ConsumerGroupCommandOptions, ConsumerGroupService} -import org.apache.kafka.clients.admin._ -import org.apache.kafka.clients.consumer.{OffsetAndMetadata, RangeAssignor} -import org.apache.kafka.common.{ConsumerGroupState, KafkaFuture, Node, TopicPartition, TopicPartitionInfo} -import org.junit.jupiter.api.Assertions.{assertEquals, assertTrue} -import org.junit.jupiter.api.Test -import org.mockito.ArgumentMatchers -import org.mockito.ArgumentMatchers._ -import org.mockito.Mockito._ -import org.mockito.ArgumentMatcher - -import scala.jdk.CollectionConverters._ -import org.apache.kafka.common.internals.KafkaFutureImpl - -class ConsumerGroupServiceTest { - - private val group = "testGroup" - private val topics = (0 until 5).map(i => s"testTopic$i") - private val numPartitions = 10 - private val topicPartitions = topics.flatMap(topic => (0 until numPartitions).map(i => new TopicPartition(topic, i))) - private val admin = mock(classOf[Admin]) - - @Test - def testAdminRequestsForDescribeOffsets(): Unit = { - val args = Array("--bootstrap-server", "localhost:9092", "--group", group, "--describe", "--offsets") - val groupService = consumerGroupService(args) - - when(admin.describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(group)), any())) - .thenReturn(describeGroupsResult(ConsumerGroupState.STABLE)) - when(admin.listConsumerGroupOffsets(ArgumentMatchers.eq(listConsumerGroupOffsetsSpec), any())) - .thenReturn(listGroupOffsetsResult(group)) - when(admin.listOffsets(offsetsArgMatcher, any())) - .thenReturn(listOffsetsResult) - - val (state, assignments) = groupService.collectGroupOffsets(group) - assertEquals(Some("Stable"), state) - assertTrue(assignments.nonEmpty) - assertEquals(topicPartitions.size, assignments.get.size) - - verify(admin, times(1)).describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(group)), any()) - verify(admin, times(1)).listConsumerGroupOffsets(ArgumentMatchers.eq(listConsumerGroupOffsetsSpec), any()) - verify(admin, times(1)).listOffsets(offsetsArgMatcher, any()) - } - - @Test - def testAdminRequestsForDescribeNegativeOffsets(): Unit = { - val args = Array("--bootstrap-server", "localhost:9092", "--group", group, "--describe", "--offsets") - val groupService = consumerGroupService(args) - - val testTopicPartition0 = new TopicPartition("testTopic1", 0); - val testTopicPartition1 = new TopicPartition("testTopic1", 1); - val testTopicPartition2 = new TopicPartition("testTopic1", 2); - val testTopicPartition3 = new TopicPartition("testTopic2", 0); - val testTopicPartition4 = new TopicPartition("testTopic2", 1); - val testTopicPartition5 = new TopicPartition("testTopic2", 2); - - // Some topic's partitions gets valid OffsetAndMetadata values, other gets nulls values (negative integers) and others aren't defined - val committedOffsets = Map( - testTopicPartition1 -> new OffsetAndMetadata(100), - testTopicPartition2 -> null, - testTopicPartition3 -> new OffsetAndMetadata(100), - testTopicPartition4 -> new OffsetAndMetadata(100), - testTopicPartition5 -> null, - ).asJava - - val resultInfo = new ListOffsetsResult.ListOffsetsResultInfo(100, System.currentTimeMillis, Optional.of(1)) - val endOffsets = Map( - testTopicPartition0 -> KafkaFuture.completedFuture(resultInfo), - testTopicPartition1 -> KafkaFuture.completedFuture(resultInfo), - testTopicPartition2 -> KafkaFuture.completedFuture(resultInfo), - testTopicPartition3 -> KafkaFuture.completedFuture(resultInfo), - testTopicPartition4 -> KafkaFuture.completedFuture(resultInfo), - testTopicPartition5 -> KafkaFuture.completedFuture(resultInfo), - ) - val assignedTopicPartitions = Set(testTopicPartition0, testTopicPartition1, testTopicPartition2) - val unassignedTopicPartitions = Set(testTopicPartition3, testTopicPartition4, testTopicPartition5) - - val consumerGroupDescription = new ConsumerGroupDescription(group, - true, - Collections.singleton(new MemberDescription("member1", Optional.of("instance1"), "client1", "host1", new MemberAssignment(assignedTopicPartitions.asJava))), - classOf[RangeAssignor].getName, - ConsumerGroupState.STABLE, - new Node(1, "localhost", 9092)) - - def offsetsArgMatcher(expectedPartitions: Set[TopicPartition]): ArgumentMatcher[util.Map[TopicPartition, OffsetSpec]] = { - topicPartitionOffsets => topicPartitionOffsets != null && topicPartitionOffsets.keySet.asScala.equals(expectedPartitions) - } - - val future = new KafkaFutureImpl[ConsumerGroupDescription]() - future.complete(consumerGroupDescription) - when(admin.describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(group)), any())) - .thenReturn(new DescribeConsumerGroupsResult(Collections.singletonMap(group, future))) - when(admin.listConsumerGroupOffsets(ArgumentMatchers.eq(listConsumerGroupOffsetsSpec), any())) - .thenReturn( - AdminClientTestUtils.listConsumerGroupOffsetsResult( - Collections.singletonMap(group, committedOffsets))) - when(admin.listOffsets( - ArgumentMatchers.argThat(offsetsArgMatcher(assignedTopicPartitions)), - any() - )).thenReturn(new ListOffsetsResult(endOffsets.filter { case (tp, _) => assignedTopicPartitions.contains(tp) }.asJava)) - when(admin.listOffsets( - ArgumentMatchers.argThat(offsetsArgMatcher(unassignedTopicPartitions)), - any() - )).thenReturn(new ListOffsetsResult(endOffsets.filter { case (tp, _) => unassignedTopicPartitions.contains(tp) }.asJava)) - - val (state, assignments) = groupService.collectGroupOffsets(group) - val returnedOffsets = assignments.map { results => - results.map { assignment => - new TopicPartition(assignment.topic.get, assignment.partition.get) -> assignment.offset - }.toMap - }.getOrElse(Map.empty) - - val expectedOffsets = Map( - testTopicPartition0 -> None, - testTopicPartition1 -> Some(100), - testTopicPartition2 -> None, - testTopicPartition3 -> Some(100), - testTopicPartition4 -> Some(100), - testTopicPartition5 -> None - ) - assertEquals(Some("Stable"), state) - assertEquals(expectedOffsets, returnedOffsets) - - verify(admin, times(1)).describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(group)), any()) - verify(admin, times(1)).listConsumerGroupOffsets(ArgumentMatchers.eq(listConsumerGroupOffsetsSpec), any()) - verify(admin, times(1)).listOffsets(ArgumentMatchers.argThat(offsetsArgMatcher(assignedTopicPartitions)), any()) - verify(admin, times(1)).listOffsets(ArgumentMatchers.argThat(offsetsArgMatcher(unassignedTopicPartitions)), any()) - } - - @Test - def testAdminRequestsForResetOffsets(): Unit = { - val args = Seq("--bootstrap-server", "localhost:9092", "--group", group, "--reset-offsets", "--to-latest") - val topicsWithoutPartitionsSpecified = topics.tail - val topicArgs = Seq("--topic", s"${topics.head}:${(0 until numPartitions).mkString(",")}") ++ - topicsWithoutPartitionsSpecified.flatMap(topic => Seq("--topic", topic)) - val groupService = consumerGroupService((args ++ topicArgs).toArray) - - when(admin.describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(group)), any())) - .thenReturn(describeGroupsResult(ConsumerGroupState.DEAD)) - when(admin.describeTopics(ArgumentMatchers.eq(topicsWithoutPartitionsSpecified.asJava), any())) - .thenReturn(describeTopicsResult(topicsWithoutPartitionsSpecified)) - when(admin.listOffsets(offsetsArgMatcher, any())) - .thenReturn(listOffsetsResult) - - val resetResult = groupService.resetOffsets() - assertEquals(Set(group), resetResult.keySet) - assertEquals(topicPartitions.toSet, resetResult(group).keySet) - - verify(admin, times(1)).describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(group)), any()) - verify(admin, times(1)).describeTopics(ArgumentMatchers.eq(topicsWithoutPartitionsSpecified.asJava), any()) - verify(admin, times(1)).listOffsets(offsetsArgMatcher, any()) - } - - private def consumerGroupService(args: Array[String]): ConsumerGroupService = { - new ConsumerGroupService(new ConsumerGroupCommandOptions(args)) { - override protected def createAdminClient(configOverrides: collection.Map[String, String]): Admin = { - admin - } - } - } - - private def describeGroupsResult(groupState: ConsumerGroupState): DescribeConsumerGroupsResult = { - val member1 = new MemberDescription("member1", Optional.of("instance1"), "client1", "host1", null) - val description = new ConsumerGroupDescription(group, - true, - Collections.singleton(member1), - classOf[RangeAssignor].getName, - groupState, - new Node(1, "localhost", 9092)) - val future = new KafkaFutureImpl[ConsumerGroupDescription]() - future.complete(description) - new DescribeConsumerGroupsResult(Collections.singletonMap(group, future)) - } - - private def listGroupOffsetsResult(groupId: String): ListConsumerGroupOffsetsResult = { - val offsets = topicPartitions.map(_ -> new OffsetAndMetadata(100)).toMap.asJava - AdminClientTestUtils.listConsumerGroupOffsetsResult(Map(groupId -> offsets).asJava) - } - - private def offsetsArgMatcher: util.Map[TopicPartition, OffsetSpec] = { - val expectedOffsets = topicPartitions.map(tp => tp -> OffsetSpec.latest).toMap - ArgumentMatchers.argThat[util.Map[TopicPartition, OffsetSpec]] { map => - map.keySet.asScala == expectedOffsets.keySet && map.values.asScala.forall(_.isInstanceOf[OffsetSpec.LatestSpec]) - } - } - - private def listOffsetsResult: ListOffsetsResult = { - val resultInfo = new ListOffsetsResult.ListOffsetsResultInfo(100, System.currentTimeMillis, Optional.of(1)) - val futures = topicPartitions.map(_ -> KafkaFuture.completedFuture(resultInfo)).toMap - new ListOffsetsResult(futures.asJava) - } - - private def describeTopicsResult(topics: Seq[String]): DescribeTopicsResult = { - val topicDescriptions = topics.map { topic => - val partitions = (0 until numPartitions).map(i => new TopicPartitionInfo(i, null, Collections.emptyList[Node], Collections.emptyList[Node])) - topic -> new TopicDescription(topic, false, partitions.asJava) - }.toMap - AdminClientTestUtils.describeTopicsResult(topicDescriptions.asJava) - } - - private def listConsumerGroupOffsetsSpec: util.Map[String, ListConsumerGroupOffsetsSpec] = { - Collections.singletonMap(group, new ListConsumerGroupOffsetsSpec()) - } -} diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandTest.java deleted file mode 100644 index 6c3712a5d218e..0000000000000 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommandTest.java +++ /dev/null @@ -1,312 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.apache.kafka.tools.consumer.group; - -import kafka.admin.ConsumerGroupCommand; -import kafka.server.KafkaConfig; -import kafka.utils.TestUtils; -import org.apache.kafka.clients.admin.AdminClientConfig; -import org.apache.kafka.clients.consumer.Consumer; -import org.apache.kafka.clients.consumer.KafkaConsumer; -import org.apache.kafka.clients.consumer.RangeAssignor; -import org.apache.kafka.common.TopicPartition; -import org.apache.kafka.common.errors.WakeupException; -import org.apache.kafka.common.serialization.StringDeserializer; -import org.junit.jupiter.api.AfterEach; -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.TestInfo; -import scala.collection.JavaConverters; -import scala.collection.Seq; - -import java.time.Duration; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.Collection; -import java.util.Collections; -import java.util.List; -import java.util.Map; -import java.util.Optional; -import java.util.Properties; -import java.util.Set; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; -import java.util.stream.Collectors; -import java.util.stream.IntStream; - -public class ConsumerGroupCommandTest extends kafka.integration.KafkaServerTestHarness { - public static final String TOPIC = "foo"; - public static final String GROUP = "test.group"; - - List consumerGroupService = new ArrayList<>(); - List consumerGroupExecutors = new ArrayList<>(); - - @Override - public Seq generateConfigs() { - List cfgs = new ArrayList<>(); - - TestUtils.createBrokerConfigs( - 1, - zkConnectOrNull(), - false, - true, - scala.None$.empty(), - scala.None$.empty(), - scala.None$.empty(), - true, - false, - false, - false, - scala.collection.immutable.Map$.MODULE$.empty(), - 1, - false, - 1, - (short) 1, - 0, - false - ).foreach(props -> { - cfgs.add(KafkaConfig.fromProps(props)); - return null; - }); - - return seq(cfgs); - } - - @BeforeEach - @Override - public void setUp(TestInfo testInfo) { - super.setUp(testInfo); - createTopic(TOPIC, 1, 1, new Properties(), listenerName(), new Properties()); - } - - @AfterEach - @Override - public void tearDown() { - consumerGroupService.forEach(ConsumerGroupCommand.ConsumerGroupService::close); - consumerGroupExecutors.forEach(AbstractConsumerGroupExecutor::shutdown); - super.tearDown(); - } - - Map committedOffsets(String topic, String group) { - try (Consumer consumer = createNoAutoCommitConsumer(group)) { - Set partitions = consumer.partitionsFor(topic).stream() - .map(partitionInfo -> new TopicPartition(partitionInfo.topic(), partitionInfo.partition())) - .collect(Collectors.toSet()); - return consumer.committed(partitions).entrySet().stream() - .filter(e -> e.getValue() != null) - .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().offset())); - } - } - - Consumer createNoAutoCommitConsumer(String group) { - Properties props = new Properties(); - props.put("bootstrap.servers", bootstrapServers(listenerName())); - props.put("group.id", group); - props.put("enable.auto.commit", "false"); - return new KafkaConsumer<>(props, new StringDeserializer(), new StringDeserializer()); - } - - ConsumerGroupCommand.ConsumerGroupService getConsumerGroupService(String[] args) { - ConsumerGroupCommand.ConsumerGroupCommandOptions opts = new ConsumerGroupCommand.ConsumerGroupCommandOptions(args); - ConsumerGroupCommand.ConsumerGroupService service = new ConsumerGroupCommand.ConsumerGroupService( - opts, - JavaConverters.asScala(Collections.singletonMap(AdminClientConfig.RETRIES_CONFIG, Integer.toString(Integer.MAX_VALUE))) - ); - - consumerGroupService.add(0, service); - return service; - } - - ConsumerGroupExecutor addConsumerGroupExecutor(int numConsumers) { - return addConsumerGroupExecutor(numConsumers, TOPIC, GROUP, RangeAssignor.class.getName(), Optional.empty(), false); - } - - ConsumerGroupExecutor addConsumerGroupExecutor(int numConsumers, String topic, String group) { - return addConsumerGroupExecutor(numConsumers, topic, group, RangeAssignor.class.getName(), Optional.empty(), false); - } - - ConsumerGroupExecutor addConsumerGroupExecutor(int numConsumers, String topic, String group, String strategy, - Optional customPropsOpt, boolean syncCommit) { - ConsumerGroupExecutor executor = new ConsumerGroupExecutor(bootstrapServers(listenerName()), numConsumers, group, - topic, strategy, customPropsOpt, syncCommit); - addExecutor(executor); - return executor; - } - - SimpleConsumerGroupExecutor addSimpleGroupExecutor(String group) { - return addSimpleGroupExecutor(Arrays.asList(new TopicPartition(TOPIC, 0)), group); - } - - SimpleConsumerGroupExecutor addSimpleGroupExecutor(Collection partitions, String group) { - SimpleConsumerGroupExecutor executor = new SimpleConsumerGroupExecutor(bootstrapServers(listenerName()), group, partitions); - addExecutor(executor); - return executor; - } - - private AbstractConsumerGroupExecutor addExecutor(AbstractConsumerGroupExecutor executor) { - consumerGroupExecutors.add(0, executor); - return executor; - } - - abstract class AbstractConsumerRunnable implements Runnable { - final String broker; - final String groupId; - final Optional customPropsOpt; - final boolean syncCommit; - - final Properties props = new Properties(); - KafkaConsumer consumer; - - boolean configured = false; - - public AbstractConsumerRunnable(String broker, String groupId, Optional customPropsOpt, boolean syncCommit) { - this.broker = broker; - this.groupId = groupId; - this.customPropsOpt = customPropsOpt; - this.syncCommit = syncCommit; - } - - void configure() { - configured = true; - configure(props); - customPropsOpt.ifPresent(props::putAll); - consumer = new KafkaConsumer<>(props); - } - - void configure(Properties props) { - props.put("bootstrap.servers", broker); - props.put("group.id", groupId); - props.put("key.deserializer", StringDeserializer.class.getName()); - props.put("value.deserializer", StringDeserializer.class.getName()); - } - - abstract void subscribe(); - - @Override - public void run() { - assert configured : "Must call configure before use"; - try { - subscribe(); - while (true) { - consumer.poll(Duration.ofMillis(Long.MAX_VALUE)); - if (syncCommit) - consumer.commitSync(); - } - } catch (WakeupException e) { - // OK - } finally { - consumer.close(); - } - } - - void shutdown() { - consumer.wakeup(); - } - } - - class ConsumerRunnable extends AbstractConsumerRunnable { - final String topic; - final String strategy; - - public ConsumerRunnable(String broker, String groupId, String topic, String strategy, - Optional customPropsOpt, boolean syncCommit) { - super(broker, groupId, customPropsOpt, syncCommit); - - this.topic = topic; - this.strategy = strategy; - } - - @Override - void configure(Properties props) { - super.configure(props); - props.put("partition.assignment.strategy", strategy); - } - - @Override - void subscribe() { - consumer.subscribe(Collections.singleton(topic)); - } - } - - class SimpleConsumerRunnable extends AbstractConsumerRunnable { - final Collection partitions; - - public SimpleConsumerRunnable(String broker, String groupId, Collection partitions) { - super(broker, groupId, Optional.empty(), false); - - this.partitions = partitions; - } - - @Override - void subscribe() { - consumer.assign(partitions); - } - } - - class AbstractConsumerGroupExecutor { - final int numThreads; - final ExecutorService executor; - final List consumers = new ArrayList<>(); - - public AbstractConsumerGroupExecutor(int numThreads) { - this.numThreads = numThreads; - this.executor = Executors.newFixedThreadPool(numThreads); - } - - void submit(AbstractConsumerRunnable consumerThread) { - consumers.add(consumerThread); - executor.submit(consumerThread); - } - - void shutdown() { - consumers.forEach(AbstractConsumerRunnable::shutdown); - executor.shutdown(); - try { - executor.awaitTermination(5000, TimeUnit.MILLISECONDS); - } catch (InterruptedException e) { - throw new RuntimeException(e); - } - } - } - - class ConsumerGroupExecutor extends AbstractConsumerGroupExecutor { - public ConsumerGroupExecutor(String broker, int numConsumers, String groupId, String topic, String strategy, - Optional customPropsOpt, boolean syncCommit) { - super(numConsumers); - IntStream.rangeClosed(1, numConsumers).forEach(i -> { - ConsumerRunnable th = new ConsumerRunnable(broker, groupId, topic, strategy, customPropsOpt, syncCommit); - th.configure(); - submit(th); - }); - } - } - - class SimpleConsumerGroupExecutor extends AbstractConsumerGroupExecutor { - public SimpleConsumerGroupExecutor(String broker, String groupId, Collection partitions) { - super(1); - - SimpleConsumerRunnable th = new SimpleConsumerRunnable(broker, groupId, partitions); - th.configure(); - submit(th); - } - } - - @SuppressWarnings({"deprecation"}) - static Seq seq(Collection seq) { - return JavaConverters.asScalaIteratorConverter(seq.iterator()).asScala().toSeq(); - } -} \ No newline at end of file From dcda50b0e2cba0542c70bd6ff61a4bc16321c22a Mon Sep 17 00:00:00 2001 From: "n.izhikov" Date: Tue, 23 Jan 2024 18:19:54 +0300 Subject: [PATCH 12/18] KAFKA-14589 ConsumerGroupServiceTest rewritten in java --- checkstyle/import-control.xml | 1 + .../group/ConsumerGroupServiceTest.java | 32 +++++++++++-------- 2 files changed, 20 insertions(+), 13 deletions(-) diff --git a/checkstyle/import-control.xml b/checkstyle/import-control.xml index b18429a620583..b6fa73b68caf6 100644 --- a/checkstyle/import-control.xml +++ b/checkstyle/import-control.xml @@ -333,6 +333,7 @@ + diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java index 7cc46acff118c..07f6f762fac8b 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java @@ -43,6 +43,7 @@ import org.mockito.ArgumentMatcher; import org.mockito.ArgumentMatchers; import scala.Option; +import scala.Some; import scala.Tuple2; import scala.collection.JavaConverters; import scala.collection.Seq; @@ -81,10 +82,10 @@ public class ConsumerGroupServiceTest { .flatMap(topic -> IntStream.range(0, NUM_PARTITIONS).mapToObj(i -> new TopicPartition(topic, i))) .collect(Collectors.toList()); - private Admin admin = mock(Admin.class); + private final Admin admin = mock(Admin.class); @Test - public void testAdminRequestsForDescribeOffsets() throws Exception { + public void testAdminRequestsForDescribeOffsets() { String[] args = new String[]{"--bootstrap-server", "localhost:9092", "--group", GROUP, "--describe", "--offsets"}; ConsumerGroupCommand.ConsumerGroupService groupService = consumerGroupService(args); @@ -96,7 +97,7 @@ public void testAdminRequestsForDescribeOffsets() throws Exception { .thenReturn(listOffsetsResult()); Tuple2, Option>> res = groupService.collectGroupOffsets(GROUP); - assertEquals(Optional.of("Stable"), res._1); + assertEquals(Some.apply("Stable"), res._1); assertTrue(res._2.isDefined()); assertEquals(TOPIC_PARTITIONS.size(), res._2.get().size()); @@ -106,7 +107,7 @@ public void testAdminRequestsForDescribeOffsets() throws Exception { } @Test - public void testAdminRequestsForDescribeNegativeOffsets() throws Exception { + public void testAdminRequestsForDescribeNegativeOffsets() { String[] args = new String[]{"--bootstrap-server", "localhost:9092", "--group", GROUP, "--describe", "--offsets"}; ConsumerGroupCommand.ConsumerGroupService groupService = consumerGroupService(args); @@ -172,11 +173,16 @@ public void testAdminRequestsForDescribeNegativeOffsets() throws Exception { Option state = res._1; Option> assignments = res._2; - Map> returnedOffsets = assignments.map(results -> - results.stream().collect(Collectors.toMap( - assignment -> new TopicPartition(assignment.topic.get(), assignment.partition.get()), - assignment -> assignment.offset)) - ).orElse(Collections.emptyMap()); + Map> returnedOffsets = new HashMap<>(); + assignments.foreach(results -> { + results.foreach(assignment -> { + returnedOffsets.put( + new TopicPartition(assignment.topic().get(), (Integer) assignment.partition().get()), + assignment.offset().isDefined() ? Optional.of((Long) assignment.offset().get()) : Optional.empty()); + return null; + }); + return null; + }); Map> expectedOffsets = new HashMap<>(); @@ -187,7 +193,7 @@ public void testAdminRequestsForDescribeNegativeOffsets() throws Exception { expectedOffsets.put(testTopicPartition4, Optional.of(100L)); expectedOffsets.put(testTopicPartition5, Optional.empty()); - assertEquals(Optional.of("Stable"), state); + assertEquals(Some.apply("Stable"), state); assertEquals(expectedOffsets, returnedOffsets); verify(admin, times(1)).describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(GROUP)), any()); @@ -214,8 +220,8 @@ public void testAdminRequestsForResetOffsets() { .thenReturn(listOffsetsResult()); scala.collection.Map> resetResult = groupService.resetOffsets(); - assertEquals(Collections.singleton(GROUP), resetResult.keySet()); - assertEquals(new HashSet<>(TOPIC_PARTITIONS), resetResult.get(GROUP).get().keys()); + assertEquals(Collections.singleton(GROUP), JavaConverters.asJava(resetResult.keySet())); + assertEquals(new HashSet<>(TOPIC_PARTITIONS), JavaConverters.asJava(resetResult.get(GROUP).get().keys().toSet())); verify(admin, times(1)).describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(GROUP)), any()); verify(admin, times(1)).describeTopics(ArgumentMatchers.eq(topicsWithoutPartitionsSpecified), any()); @@ -223,7 +229,7 @@ public void testAdminRequestsForResetOffsets() { } private ConsumerGroupCommand.ConsumerGroupService consumerGroupService(String[] args) { - return new ConsumerGroupCommand.ConsumerGroupService(new ConsumerGroupCommandOptions(args), JavaConverters.asScala(Collections.emptyMap())) { + return new ConsumerGroupCommand.ConsumerGroupService(new kafka.admin.ConsumerGroupCommand.ConsumerGroupCommandOptions(args), JavaConverters.asScala(Collections.emptyMap())) { @Override public Admin createAdminClient(scala.collection.Map configOverrides) { return admin; From 2cbb2414b35ae2be244aabcd6b4df41cafd88695 Mon Sep 17 00:00:00 2001 From: "n.izhikov" Date: Tue, 23 Jan 2024 18:20:50 +0300 Subject: [PATCH 13/18] KAFKA-14589 ConsumerGroupServiceTest rewritten in java --- .../apache/kafka/tools/reassign/ReassignPartitionsCommand.java | 1 - .../main/java/org/apache/kafka/tools/{ => reassign}/Tuple2.java | 2 +- .../kafka/tools/reassign/ReassignPartitionsIntegrationTest.java | 1 - .../apache/kafka/tools/reassign/ReassignPartitionsUnitTest.java | 1 - 4 files changed, 1 insertion(+), 4 deletions(-) rename tools/src/main/java/org/apache/kafka/tools/{ => reassign}/Tuple2.java (97%) diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java b/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java index 3623346980743..6a2d9b9346863 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java @@ -49,7 +49,6 @@ import org.apache.kafka.server.util.json.JsonValue; import org.apache.kafka.tools.TerseException; import org.apache.kafka.tools.ToolsUtils; -import org.apache.kafka.tools.Tuple2; import java.io.IOException; import java.util.ArrayList; diff --git a/tools/src/main/java/org/apache/kafka/tools/Tuple2.java b/tools/src/main/java/org/apache/kafka/tools/reassign/Tuple2.java similarity index 97% rename from tools/src/main/java/org/apache/kafka/tools/Tuple2.java rename to tools/src/main/java/org/apache/kafka/tools/reassign/Tuple2.java index 02dd4bf3981a8..a9e84317cf718 100644 --- a/tools/src/main/java/org/apache/kafka/tools/Tuple2.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/Tuple2.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.tools; +package org.apache.kafka.tools.reassign; import java.util.Objects; diff --git a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsIntegrationTest.java b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsIntegrationTest.java index 83fc665e3e487..b71331a6b233f 100644 --- a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsIntegrationTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsIntegrationTest.java @@ -43,7 +43,6 @@ import org.apache.kafka.common.utils.Time; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.tools.TerseException; -import org.apache.kafka.tools.Tuple2; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Timeout; import org.junit.jupiter.params.ParameterizedTest; diff --git a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsUnitTest.java b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsUnitTest.java index 699e048fb206e..d0e1cb4baaf69 100644 --- a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsUnitTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsUnitTest.java @@ -31,7 +31,6 @@ import org.apache.kafka.common.utils.Time; import org.apache.kafka.server.common.AdminCommandFailedException; import org.apache.kafka.server.common.AdminOperationException; -import org.apache.kafka.tools.Tuple2; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; From 8c32518ea5cec448b0903085ebfa52beab3d8cb9 Mon Sep 17 00:00:00 2001 From: "n.izhikov" Date: Tue, 23 Jan 2024 18:50:33 +0300 Subject: [PATCH 14/18] KAFKA-14589 ConsumerGroupServiceTest rewritten in java --- .../consumer/group/ConsumerGroupServiceTest.java | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java index 07f6f762fac8b..13a463a0ad095 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java @@ -47,6 +47,7 @@ import scala.Tuple2; import scala.collection.JavaConverters; import scala.collection.Seq; +import scala.collection.immutable.Map$; import java.util.ArrayList; import java.util.Arrays; @@ -220,8 +221,8 @@ public void testAdminRequestsForResetOffsets() { .thenReturn(listOffsetsResult()); scala.collection.Map> resetResult = groupService.resetOffsets(); - assertEquals(Collections.singleton(GROUP), JavaConverters.asJava(resetResult.keySet())); - assertEquals(new HashSet<>(TOPIC_PARTITIONS), JavaConverters.asJava(resetResult.get(GROUP).get().keys().toSet())); + assertEquals(set(Collections.singletonList(GROUP)), resetResult.keySet()); + assertEquals(set(TOPIC_PARTITIONS), resetResult.get(GROUP).get().keys().toSet()); verify(admin, times(1)).describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(GROUP)), any()); verify(admin, times(1)).describeTopics(ArgumentMatchers.eq(topicsWithoutPartitionsSpecified), any()); @@ -229,7 +230,7 @@ public void testAdminRequestsForResetOffsets() { } private ConsumerGroupCommand.ConsumerGroupService consumerGroupService(String[] args) { - return new ConsumerGroupCommand.ConsumerGroupService(new kafka.admin.ConsumerGroupCommand.ConsumerGroupCommandOptions(args), JavaConverters.asScala(Collections.emptyMap())) { + return new ConsumerGroupCommand.ConsumerGroupService(new kafka.admin.ConsumerGroupCommand.ConsumerGroupCommandOptions(args), Map$.MODULE$.empty()) { @Override public Admin createAdminClient(scala.collection.Map configOverrides) { return admin; @@ -290,4 +291,9 @@ private DescribeTopicsResult describeTopicsResult(Collection topics) { private Map listConsumerGroupOffsetsSpec() { return Collections.singletonMap(GROUP, new ListConsumerGroupOffsetsSpec()); } + + @SuppressWarnings({"deprecation"}) + private static scala.collection.immutable.Set set(final Collection set) { + return JavaConverters.asScalaSet(new HashSet<>(set)).toSet(); + } } From cdcea2afefd9da6252960d748237477777d77744 Mon Sep 17 00:00:00 2001 From: "n.izhikov" Date: Wed, 24 Jan 2024 20:51:36 +0300 Subject: [PATCH 15/18] KAFKA-14589 Review fixes --- .../consumer/group/ConsumerGroupServiceTest.java | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java index 13a463a0ad095..3611406d23e81 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java @@ -97,10 +97,10 @@ public void testAdminRequestsForDescribeOffsets() { when(admin.listOffsets(offsetsArgMatcher(), any())) .thenReturn(listOffsetsResult()); - Tuple2, Option>> res = groupService.collectGroupOffsets(GROUP); - assertEquals(Some.apply("Stable"), res._1); - assertTrue(res._2.isDefined()); - assertEquals(TOPIC_PARTITIONS.size(), res._2.get().size()); + Tuple2, Option>> statesAndAssignments = groupService.collectGroupOffsets(GROUP); + assertEquals(Some.apply("Stable"), statesAndAssignments._1); + assertTrue(statesAndAssignments._2.isDefined()); + assertEquals(TOPIC_PARTITIONS.size(), statesAndAssignments._2.get().size()); verify(admin, times(1)).describeConsumerGroups(ArgumentMatchers.eq(Collections.singletonList(GROUP)), any()); verify(admin, times(1)).listConsumerGroupOffsets(ArgumentMatchers.eq(listConsumerGroupOffsetsSpec()), any()); @@ -119,7 +119,7 @@ public void testAdminRequestsForDescribeNegativeOffsets() { TopicPartition testTopicPartition4 = new TopicPartition("testTopic2", 1); TopicPartition testTopicPartition5 = new TopicPartition("testTopic2", 2); - // Some topic's partitions gets valid OffsetAndMetadata values, other gets nulls values (negative integers) and others aren't defined + // Some topic's partitions gets valid OffsetAndMetadata values, other gets nulls values and others aren't defined Map committedOffsets = new HashMap<>(); committedOffsets.put(testTopicPartition1, new OffsetAndMetadata(100)); @@ -170,9 +170,9 @@ public void testAdminRequestsForDescribeNegativeOffsets() { )).thenReturn(new ListOffsetsResult(endOffsets.entrySet().stream().filter(e -> unassignedTopicPartitions.contains(e.getKey())) .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)))); - Tuple2, Option>> res = groupService.collectGroupOffsets(GROUP); - Option state = res._1; - Option> assignments = res._2; + Tuple2, Option>> statesAndAssignments = groupService.collectGroupOffsets(GROUP); + Option state = statesAndAssignments._1; + Option> assignments = statesAndAssignments._2; Map> returnedOffsets = new HashMap<>(); assignments.foreach(results -> { From e10f19e7541279ab2e4df3849e6bc16701d89774 Mon Sep 17 00:00:00 2001 From: "n.izhikov" Date: Wed, 24 Jan 2024 20:59:21 +0300 Subject: [PATCH 16/18] KAFKA-14589 Returning strange comment. --- .../kafka/tools/consumer/group/ConsumerGroupServiceTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java index 3611406d23e81..34230cb607a6d 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java @@ -119,7 +119,7 @@ public void testAdminRequestsForDescribeNegativeOffsets() { TopicPartition testTopicPartition4 = new TopicPartition("testTopic2", 1); TopicPartition testTopicPartition5 = new TopicPartition("testTopic2", 2); - // Some topic's partitions gets valid OffsetAndMetadata values, other gets nulls values and others aren't defined + // Some topic's partitions gets valid OffsetAndMetadata values, other gets nulls values (negative integers) and others aren't defined Map committedOffsets = new HashMap<>(); committedOffsets.put(testTopicPartition1, new OffsetAndMetadata(100)); From ac3bb20a6f6ea6c9a95e6c3d77afd1bf46406bba Mon Sep 17 00:00:00 2001 From: "n.izhikov" Date: Thu, 25 Jan 2024 00:40:52 +0300 Subject: [PATCH 17/18] KAFKA-14589 Returning strange comment. --- .../tools/consumer/group/ConsumerGroupServiceTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java index 34230cb607a6d..c5acabb6f5dbc 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java @@ -254,14 +254,14 @@ private DescribeConsumerGroupsResult describeGroupsResult(ConsumerGroupState gro private ListConsumerGroupOffsetsResult listGroupOffsetsResult(String groupId) { Map offsets = TOPIC_PARTITIONS.stream().collect(Collectors.toMap( Function.identity(), - tp -> new OffsetAndMetadata(100))); + unused -> new OffsetAndMetadata(100))); return AdminClientTestUtils.listConsumerGroupOffsetsResult(Collections.singletonMap(groupId, offsets)); } private Map offsetsArgMatcher() { Map expectedOffsets = TOPIC_PARTITIONS.stream().collect(Collectors.toMap( Function.identity(), - tp -> OffsetSpec.latest() + unused -> OffsetSpec.latest() )); return ArgumentMatchers.argThat(map -> Objects.equals(map.keySet(), expectedOffsets.keySet()) && map.values().stream().allMatch(v -> v instanceof OffsetSpec.LatestSpec) @@ -272,7 +272,7 @@ private ListOffsetsResult listOffsetsResult() { ListOffsetsResultInfo resultInfo = new ListOffsetsResultInfo(100, System.currentTimeMillis(), Optional.of(1)); Map> futures = TOPIC_PARTITIONS.stream().collect(Collectors.toMap( Function.identity(), - tp -> KafkaFuture.completedFuture(resultInfo))); + unused -> KafkaFuture.completedFuture(resultInfo))); return new ListOffsetsResult(futures); } From 2a39d6123921534c7f6404c6bff315b4b5f60361 Mon Sep 17 00:00:00 2001 From: "n.izhikov" Date: Thu, 25 Jan 2024 00:55:09 +0300 Subject: [PATCH 18/18] KAFKA-14589 Returning strange comment. --- .../tools/consumer/group/ConsumerGroupServiceTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java index c5acabb6f5dbc..8c75823648f83 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ConsumerGroupServiceTest.java @@ -254,14 +254,14 @@ private DescribeConsumerGroupsResult describeGroupsResult(ConsumerGroupState gro private ListConsumerGroupOffsetsResult listGroupOffsetsResult(String groupId) { Map offsets = TOPIC_PARTITIONS.stream().collect(Collectors.toMap( Function.identity(), - unused -> new OffsetAndMetadata(100))); + __ -> new OffsetAndMetadata(100))); return AdminClientTestUtils.listConsumerGroupOffsetsResult(Collections.singletonMap(groupId, offsets)); } private Map offsetsArgMatcher() { Map expectedOffsets = TOPIC_PARTITIONS.stream().collect(Collectors.toMap( Function.identity(), - unused -> OffsetSpec.latest() + __ -> OffsetSpec.latest() )); return ArgumentMatchers.argThat(map -> Objects.equals(map.keySet(), expectedOffsets.keySet()) && map.values().stream().allMatch(v -> v instanceof OffsetSpec.LatestSpec) @@ -272,7 +272,7 @@ private ListOffsetsResult listOffsetsResult() { ListOffsetsResultInfo resultInfo = new ListOffsetsResultInfo(100, System.currentTimeMillis(), Optional.of(1)); Map> futures = TOPIC_PARTITIONS.stream().collect(Collectors.toMap( Function.identity(), - unused -> KafkaFuture.completedFuture(resultInfo))); + __ -> KafkaFuture.completedFuture(resultInfo))); return new ListOffsetsResult(futures); }