This is an automated email from the ASF dual-hosted git repository.
AndrewJSchofield pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 33b93fc4114 KAFKA-20524: CSV reset offset plan for
kafka-share-groups.sh (#22197)
33b93fc4114 is described below
commit 33b93fc411441c95a89f40e17006f2b1ae2c77da
Author: Andrew Schofield <[email protected]>
AuthorDate: Tue Jun 16 15:30:51 2026 +0100
KAFKA-20524: CSV reset offset plan for kafka-share-groups.sh (#22197)
This is the implementation of
https://cwiki.apache.org/confluence/display/KAFKA/KIP-1323%3A+Initialization+of+share+group+offsets+from+a+specific+offset+or+a+file.
Reviewers: Chia-Ping Tsai <[email protected]>, Apoorv Mittal
<[email protected]>
---
.../apache/kafka/tools/GroupOffsetsResetter.java | 64 ++++---
.../group/ConsumerGroupCommandOptions.java | 4 +-
.../tools/consumer/group/ShareGroupCommand.java | 119 ++++++++++--
.../consumer/group/ShareGroupCommandOptions.java | 63 ++++--
.../consumer/group/ShareGroupCommandTest.java | 213 ++++++++++++++++++++-
5 files changed, 387 insertions(+), 76 deletions(-)
diff --git
a/tools/src/main/java/org/apache/kafka/tools/GroupOffsetsResetter.java
b/tools/src/main/java/org/apache/kafka/tools/GroupOffsetsResetter.java
index 16725f9fc83..9c7afe92bb9 100644
--- a/tools/src/main/java/org/apache/kafka/tools/GroupOffsetsResetter.java
+++ b/tools/src/main/java/org/apache/kafka/tools/GroupOffsetsResetter.java
@@ -22,6 +22,7 @@ import org.apache.kafka.clients.admin.DescribeTopicsOptions;
import org.apache.kafka.clients.admin.ListOffsetsOptions;
import org.apache.kafka.clients.admin.ListOffsetsResult;
import org.apache.kafka.clients.admin.OffsetSpec;
+import org.apache.kafka.clients.admin.SharePartitionOffsetInfo;
import org.apache.kafka.clients.admin.TopicDescription;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
@@ -117,7 +118,6 @@ public class GroupOffsetsResetter {
isOldCsvFormat = true;
}
} catch (IOException e) {
- e.printStackTrace();
// Ignore.
}
@@ -159,14 +159,14 @@ public class GroupOffsetsResetter {
if (logEndOffset != null) {
if (logEndOffset instanceof LogOffset && offset > ((LogOffset)
logEndOffset).value) {
long endOffset = ((LogOffset) logEndOffset).value;
- LOGGER.warn("New offset (" + offset + ") is higher than
latest offset for topic partition " + topicPartition + ". Value will be set to
" + endOffset);
+ LOGGER.warn("New offset (" + offset + ") is higher than
latest offset for topic partition " + topicPartition + ". Value will be set to
" + endOffset + ".");
res.put(topicPartition, endOffset);
} else {
LogOffsetResult logStartOffset =
logStartOffsets.get(topicPartition);
if (logStartOffset instanceof LogOffset && offset <
((LogOffset) logStartOffset).value) {
long startOffset = ((LogOffset) logStartOffset).value;
- LOGGER.warn("New offset (" + offset + ") is lower than
earliest offset for topic partition " + topicPartition + ". Value will be set
to " + startOffset);
+ LOGGER.warn("New offset (" + offset + ") is lower than
earliest offset for topic partition " + topicPartition + ". Value will be set
to " + startOffset + ".");
res.put(topicPartition, startOffset);
} else
res.put(topicPartition, offset);
@@ -390,7 +390,7 @@ public class GroupOffsetsResetter {
Map<TopicPartition, OffsetAndMetadata> resetPlanForGroup =
resetPlan.get(groupId);
if (resetPlanForGroup == null) {
- printError("No reset plan for group " + groupId + " found",
Optional.empty());
+ printError("No reset plan for group " + groupId + " found.",
Optional.empty());
return Map.<TopicPartition, OffsetAndMetadata>of();
}
@@ -404,39 +404,39 @@ public class GroupOffsetsResetter {
}
public Map<TopicPartition, OffsetAndMetadata>
resetToCurrent(Collection<TopicPartition> partitionsToReset,
Map<TopicPartition, OffsetAndMetadata> currentCommittedOffsets) {
- Collection<TopicPartition> partitionsToResetWithCommittedOffset = new
ArrayList<>();
- Collection<TopicPartition> partitionsToResetWithoutCommittedOffset =
new ArrayList<>();
+ Map<Boolean, List<TopicPartition>> partitioned =
partitionsToReset.stream().collect(Collectors.partitioningBy(currentCommittedOffsets::containsKey));
+ List<TopicPartition> partitionsToResetWithStartOffset =
partitioned.get(true);
+ List<TopicPartition> partitionsToResetWithoutStartOffset =
partitioned.get(false);
- for (TopicPartition topicPartition : partitionsToReset) {
- if (currentCommittedOffsets.containsKey(topicPartition))
- partitionsToResetWithCommittedOffset.add(topicPartition);
- else
- partitionsToResetWithoutCommittedOffset.add(topicPartition);
- }
+ Map<TopicPartition, OffsetAndMetadata>
preparedOffsetsForPartitionsWithStartOffset =
partitionsToResetWithStartOffset.stream()
+ .collect(Collectors.toMap(Function.identity(), topicPartition ->
new OffsetAndMetadata(currentCommittedOffsets.get(topicPartition).offset())));
- Map<TopicPartition, OffsetAndMetadata>
preparedOffsetsForPartitionsWithCommittedOffset =
partitionsToResetWithCommittedOffset.stream()
- .collect(Collectors.toMap(Function.identity(), topicPartition -> {
- OffsetAndMetadata committedOffset =
currentCommittedOffsets.get(topicPartition);
+ getLogEndOffsets(partitionsToResetWithoutStartOffset).forEach((tp,
logOffsetResult) -> {
+ if (!(logOffsetResult instanceof GroupOffsetsResetter.LogOffset
logOffset)) {
+ throw new IllegalStateException("Error getting ending offset
of topic partition: " + tp);
+ }
+ preparedOffsetsForPartitionsWithStartOffset.put(tp, new
OffsetAndMetadata(logOffset.value));
+ });
- if (committedOffset == null) {
- throw new IllegalStateException("Expected a valid current
offset for topic partition: " + topicPartition);
- }
+ return preparedOffsetsForPartitionsWithStartOffset;
+ }
- return new OffsetAndMetadata(committedOffset.offset());
- }));
+ public Map<TopicPartition, OffsetAndMetadata>
resetToCurrentForShareGroup(Collection<TopicPartition> partitionsToReset,
Map<TopicPartition, SharePartitionOffsetInfo> currentOffsetInfo) {
+ Map<Boolean, List<TopicPartition>> partitioned =
partitionsToReset.stream().collect(Collectors.partitioningBy(currentOffsetInfo::containsKey));
+ List<TopicPartition> partitionsToResetWithStartOffset =
partitioned.get(true);
+ List<TopicPartition> partitionsToResetWithoutStartOffset =
partitioned.get(false);
- Map<TopicPartition, OffsetAndMetadata>
preparedOffsetsForPartitionsWithoutCommittedOffset =
- getLogEndOffsets(partitionsToResetWithoutCommittedOffset)
-
.entrySet().stream().collect(Collectors.toMap(Map.Entry::getKey, e -> {
- if (!(e.getValue() instanceof
GroupOffsetsResetter.LogOffset)) {
- CommandLineUtils.printUsageAndExit(parser, "Error
getting ending offset of topic partition: " + e.getKey());
- }
- return new
OffsetAndMetadata(((GroupOffsetsResetter.LogOffset) e.getValue()).value);
- }));
+ Map<TopicPartition, OffsetAndMetadata>
preparedOffsetsForPartitionsWithStartOffset =
partitionsToResetWithStartOffset.stream()
+ .collect(Collectors.toMap(Function.identity(), topicPartition ->
new OffsetAndMetadata(currentOffsetInfo.get(topicPartition).startOffset())));
-
preparedOffsetsForPartitionsWithCommittedOffset.putAll(preparedOffsetsForPartitionsWithoutCommittedOffset);
+ getLogEndOffsets(partitionsToResetWithoutStartOffset).forEach((tp,
logOffsetResult) -> {
+ if (!(logOffsetResult instanceof GroupOffsetsResetter.LogOffset
logOffset)) {
+ throw new IllegalStateException("Error getting ending offset
of topic partition: " + tp);
+ }
+ preparedOffsetsForPartitionsWithStartOffset.put(tp, new
OffsetAndMetadata(logOffset.value));
+ });
- return preparedOffsetsForPartitionsWithCommittedOffset;
+ return preparedOffsetsForPartitionsWithStartOffset;
}
public void checkAllTopicPartitionsValid(Collection<TopicPartition>
partitionsToReset) {
@@ -519,9 +519,11 @@ public class GroupOffsetsResetter {
public GroupOffsetsResetterOptions(
List<String> groupOpt,
+ List<Long> resetToOffsetOpt,
+ List<String> resetFromFileOpt,
List<String> resetToDatetimeOpt,
long timeoutMsOpt) {
- this(groupOpt, null, null, resetToDatetimeOpt, null, null,
timeoutMsOpt);
+ this(groupOpt, resetToOffsetOpt, resetFromFileOpt,
resetToDatetimeOpt, null, null, timeoutMsOpt);
}
}
}
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 94cc72aa097..4b2ecaa4d01 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
@@ -47,13 +47,13 @@ public class ConsumerGroupCommandOptions extends
CommandDefaultOptions {
private static final String COMMAND_CONFIG_DOC = "Property file containing
configs to be passed to Admin Client and Consumer.";
private 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 +
+ "Additionally, the --export option is used to export the offsets in
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'.";
private static final String DRY_RUN_DOC = "Only show results without
executing changes on Consumer Groups. Supported operations: reset-offsets.";
private static final String EXECUTE_DOC = "Execute operation. Supported
operations: reset-offsets.";
- private static final String EXPORT_DOC = "Export operation execution to a
CSV file. Supported operations: reset-offsets.";
+ private static final String EXPORT_DOC = "Export offset information in CSV
format. Supported operations: reset-offsets.";
private static final String RESET_TO_OFFSET_DOC = "Reset offsets to a
specific offset.";
private static final String RESET_FROM_FILE_DOC = "Reset offsets to values
defined in CSV file.";
private static final String RESET_TO_DATETIME_DOC = "Reset offsets to
offset from datetime. Format: 'YYYY-MM-DDThh:mm:ss.sss'";
diff --git
a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java
b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java
index aa02dd22d3b..d06ee58bbd8 100644
---
a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java
+++
b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommand.java
@@ -46,6 +46,9 @@ import org.apache.kafka.common.utils.Utils;
import org.apache.kafka.server.util.CommandLineUtils;
import org.apache.kafka.tools.GroupOffsetsResetter;
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.ObjectWriter;
+
import java.io.IOException;
import java.util.AbstractMap.SimpleImmutableEntry;
import java.util.ArrayList;
@@ -98,7 +101,13 @@ public class ShareGroupCommand {
} else if (opts.options.has(opts.deleteOpt)) {
shareGroupService.deleteShareGroups();
} else if (opts.options.has(opts.resetOffsetsOpt)) {
- shareGroupService.resetOffsets();
+ Map<String, Map<TopicPartition, OffsetAndMetadata>>
offsetsToReset = shareGroupService.resetOffsets();
+ if (opts.options.has(opts.exportOpt)) {
+ String exported =
shareGroupService.exportOffsetsToCsv(offsetsToReset);
+ System.out.println(exported);
+ } else {
+ GroupOffsetsResetter.printOffsetsToReset(offsetsToReset);
+ }
} else if (opts.options.has(opts.deleteOffsetsOpt)) {
shareGroupService.deleteOffsets();
}
@@ -150,6 +159,8 @@ public class ShareGroupCommand {
private GroupOffsetsResetter.GroupOffsetsResetterOptions
getGroupOffsetsResetterOptions(ShareGroupCommandOptions opts) {
return
new
GroupOffsetsResetter.GroupOffsetsResetterOptions(opts.options.valuesOf(opts.groupOpt),
+ opts.options.valuesOf(opts.resetToOffsetOpt),
+ opts.options.valuesOf(opts.resetFromFileOpt),
opts.options.valuesOf(opts.resetToDatetimeOpt),
opts.options.valueOf(opts.timeoutMsOpt));
}
@@ -173,8 +184,15 @@ public class ShareGroupCommand {
.timeoutMs(opts.options.valueOf(opts.timeoutMsOpt).intValue()));
Collection<GroupListing> listings = result.all().get();
return
listings.stream().map(GroupListing::groupId).collect(Collectors.toList());
- } catch (InterruptedException | ExecutionException e) {
- throw new RuntimeException(e);
+ } catch (InterruptedException ie) {
+ throw new RuntimeException(ie);
+ } catch (ExecutionException ee) {
+ Throwable cause = ee.getCause();
+ if (cause instanceof KafkaException) {
+ throw (KafkaException) cause;
+ } else {
+ throw new RuntimeException(cause);
+ }
}
}
@@ -184,8 +202,15 @@ public class ShareGroupCommand {
.timeoutMs(opts.options.valueOf(opts.timeoutMsOpt).intValue()));
Collection<GroupListing> listings = result.all().get();
return listings.stream().toList();
- } catch (InterruptedException | ExecutionException e) {
- throw new RuntimeException(e);
+ } catch (InterruptedException ie) {
+ throw new RuntimeException(ie);
+ } catch (ExecutionException ee) {
+ Throwable cause = ee.getCause();
+ if (cause instanceof KafkaException) {
+ throw (KafkaException) cause;
+ } else {
+ throw new RuntimeException(cause);
+ }
}
}
@@ -242,6 +267,26 @@ public class ShareGroupCommand {
}
}
+ String exportOffsetsToCsv(Map<String, Map<TopicPartition,
OffsetAndMetadata>> offsetsForGroups) {
+ ObjectWriter csvWriter =
CsvUtils.writerFor(CsvUtils.CsvRecordNoGroup.class);
+
+ return offsetsForGroups.entrySet().stream().flatMap(e -> {
+ Map<TopicPartition, OffsetAndMetadata> partitionInfo =
e.getValue();
+
+ return partitionInfo.entrySet().stream().map(e1 -> {
+ TopicPartition k = e1.getKey();
+ OffsetAndMetadata v = e1.getValue();
+ Object csvRecord = new
CsvUtils.CsvRecordNoGroup(k.topic(), k.partition(), v.offset());
+
+ try {
+ return csvWriter.writeValueAsString(csvRecord);
+ } catch (JsonProcessingException err) {
+ throw new RuntimeException(err);
+ }
+ });
+ }).collect(Collectors.joining());
+ }
+
Map<String, Throwable> deleteShareGroups() {
List<GroupListing> shareGroupIds = listDetailedShareGroups();
List<String> groupIds = opts.options.has(opts.allGroupsOpt)
@@ -382,32 +427,42 @@ public class ShareGroupCommand {
return new SimpleImmutableEntry<>(topLevelException,
topicLevelResult);
}
- void resetOffsets() {
+ Map<String, Map<TopicPartition, OffsetAndMetadata>> resetOffsets() {
+ Map<String, Map<TopicPartition, OffsetAndMetadata>> result = new
HashMap<>();
+
String groupId = opts.options.valueOf(opts.groupOpt);
try {
ShareGroupDescription shareGroupDescription =
describeShareGroups(List.of(groupId)).get(groupId);
if
(!(GroupState.EMPTY.equals(shareGroupDescription.groupState()) ||
GroupState.DEAD.equals(shareGroupDescription.groupState()))) {
CommandLineUtils.printErrorAndExit(String.format("Share
group '%s' is not empty.", groupId));
}
- resetOffsetsForInactiveGroup(groupId);
+ Map<TopicPartition, OffsetAndMetadata> offsetsToReset =
resetOffsetsForInactiveGroup(groupId);
+ if (!offsetsToReset.isEmpty()) {
+ result.put(groupId, offsetsToReset);
+ }
} catch (InterruptedException ie) {
throw new RuntimeException(ie);
} catch (ExecutionException ee) {
Throwable cause = ee.getCause();
if (cause instanceof GroupIdNotFoundException) {
- resetOffsetsForInactiveGroup(groupId);
+ Map<TopicPartition, OffsetAndMetadata> offsetsToReset =
resetOffsetsForInactiveGroup(groupId);
+ if (!offsetsToReset.isEmpty()) {
+ result.put(groupId, offsetsToReset);
+ }
} else if (cause instanceof KafkaException) {
CommandLineUtils.printErrorAndExit(cause.getMessage());
} else {
throw new RuntimeException(cause);
}
}
+
+ return result;
}
- private void resetOffsetsForInactiveGroup(String groupId) {
+ private Map<TopicPartition, OffsetAndMetadata>
resetOffsetsForInactiveGroup(String groupId) {
try {
Collection<TopicPartition> partitionsToReset =
getPartitionsToReset(groupId);
- Map<TopicPartition, OffsetAndMetadata> offsetsToReset =
prepareOffsetsToReset(partitionsToReset);
+ Map<TopicPartition, OffsetAndMetadata> offsetsToReset =
prepareOffsetsToReset(groupId, partitionsToReset);
boolean dryRun = opts.options.has(opts.dryRunOpt) ||
!opts.options.has(opts.executeOpt);
if (!dryRun) {
adminClient.alterShareGroupOffsets(groupId,
@@ -418,7 +473,7 @@ public class ShareGroupCommand {
withTimeoutMs(new AlterShareGroupOffsetsOptions())
).all().get();
}
- GroupOffsetsResetter.printOffsetsToReset(Map.of(groupId,
offsetsToReset));
+ return offsetsToReset;
} catch (InterruptedException ie) {
throw new RuntimeException(ie);
} catch (ExecutionException ee) {
@@ -448,17 +503,42 @@ public class ShareGroupCommand {
return partitionsToReset;
}
- private Map<TopicPartition, OffsetAndMetadata>
prepareOffsetsToReset(Collection<TopicPartition> partitionsToReset) {
+ private Map<TopicPartition, SharePartitionOffsetInfo>
getOffsetInfo(String groupId) {
+ try {
+ return adminClient.listShareGroupOffsets(
+ Map.of(groupId, new ListShareGroupOffsetsSpec()),
+ withTimeoutMs(new ListShareGroupOffsetsOptions())
+ ).partitionsToOffsetInfo(groupId).get();
+ } catch (InterruptedException ie) {
+ throw new RuntimeException(ie);
+ } catch (ExecutionException ee) {
+ Throwable cause = ee.getCause();
+ if (cause instanceof KafkaException) {
+ throw (KafkaException) cause;
+ } else {
+ throw new RuntimeException(cause);
+ }
+ }
+ }
+
+ private Map<TopicPartition, OffsetAndMetadata>
prepareOffsetsToReset(String groupId, Collection<TopicPartition>
partitionsToReset) {
groupOffsetsResetter.checkAllTopicPartitionsValid(partitionsToReset);
- if (opts.options.has(opts.resetToEarliestOpt)) {
+ if (opts.options.has(opts.resetToOffsetOpt)) {
+ return groupOffsetsResetter.resetToOffset(partitionsToReset);
+ } else if (opts.options.has(opts.resetToEarliestOpt)) {
return groupOffsetsResetter.resetToEarliest(partitionsToReset);
} else if (opts.options.has(opts.resetToLatestOpt)) {
return groupOffsetsResetter.resetToLatest(partitionsToReset);
} else if (opts.options.has(opts.resetToDatetimeOpt)) {
return groupOffsetsResetter.resetToDateTime(partitionsToReset);
+ } else if (opts.options.has(opts.resetFromFileOpt)) {
+ return groupOffsetsResetter.resetFromFile(groupId);
+ } else if (opts.options.has(opts.resetToCurrentOpt)) {
+ Map<TopicPartition, SharePartitionOffsetInfo> currentOffsets =
getOffsetInfo(groupId);
+ return
groupOffsetsResetter.resetToCurrentForShareGroup(partitionsToReset,
currentOffsets);
}
CommandLineUtils
- .printUsageAndExit(opts.parser, String.format("Option '%s'
requires one of the following scenarios: %s", opts.resetOffsetsOpt,
opts.allResetOffsetScenarioOpts));
+ .printUsageAndExit(opts.parser, String.format("Option '%s'
requires one of the following scenarios: %s", opts.resetOffsetsOpt,
opts.allResetOffsetsScenarioOpts));
return null;
}
@@ -502,8 +582,15 @@ public class ShareGroupCommand {
Set<SharePartitionOffsetInformation> partitionOffsets =
mapOffsetInfoToSharePartitionInformation(groupId, offsetInfoMap);
groupOffsets.put(groupId, new
SimpleImmutableEntry<>(shareGroup, partitionOffsets));
- } catch (InterruptedException | ExecutionException e) {
- throw new RuntimeException(e);
+ } catch (InterruptedException ie) {
+ throw new RuntimeException(ie);
+ } catch (ExecutionException ee) {
+ Throwable cause = ee.getCause();
+ if (cause instanceof KafkaException) {
+ throw (KafkaException) cause;
+ } else {
+ throw new RuntimeException(cause);
+ }
}
});
diff --git
a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandOptions.java
b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandOptions.java
index b227ac71374..a3f1eff877a 100644
---
a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandOptions.java
+++
b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandOptions.java
@@ -48,14 +48,19 @@ public class ShareGroupCommandOptions extends
CommandDefaultOptions {
private static final String COMMAND_CONFIG_DOC = "Property file containing
configs to be passed to Admin Client.";
private static final String RESET_OFFSETS_DOC = "Reset offsets of share
group. Supports one share group at the time, and instances must be inactive." +
NL +
"Has 2 execution options: --dry-run to plan which offsets to reset,
and --execute to reset the offsets. " + NL +
- "You must choose one of the following reset specifications:
--to-datetime, --to-earliest, --to-latest." + NL +
- "To define the scope use --all-topics or --topic." + NL +
+ "Additionally, the --export option is used to export the offsets in
CSV format." + NL +
+ "You must choose one of the following reset specifications:
--to-datetime, --to-earliest, --to-latest, --from-file, --to-current,
--to-offset." + NL +
+ "To define the scope, use --all-topics, --topic or --from-file." + NL +
"Fails if neither '--dry-run' nor '--execute' is specified.";
private static final String DRY_RUN_DOC = "Only show results without
executing changes on share groups. Supported operations: reset-offsets.";
private static final String EXECUTE_DOC = "Execute operation. Supported
operations: reset-offsets.";
+ private static final String EXPORT_DOC = "Export offset information in CSV
format. Supported operations: reset-offsets.";
+ private static final String RESET_TO_OFFSET_DOC = "Reset offsets to a
specific offset.";
+ private static final String RESET_FROM_FILE_DOC = "Reset offsets to values
defined in CSV file.";
private static final String RESET_TO_DATETIME_DOC = "Reset offsets to
offset from datetime. Format: 'YYYY-MM-DDThh:mm:ss.sss'";
private static final String RESET_TO_EARLIEST_DOC = "Reset offsets to
earliest offset.";
private static final String RESET_TO_LATEST_DOC = "Reset offsets to latest
offset.";
+ private static final String RESET_TO_CURRENT_DOC = "Reset offsets to
current offset.";
private static final String MEMBERS_DOC = "Describe members of the group.
This option may be used with the '--describe' option only.";
private static final String OFFSETS_DOC = "Describe the group and list all
topic partitions in the group along with their offset information. " +
"This is the default sub-action and may be used with the '--describe'
option only.";
@@ -79,10 +84,14 @@ public class ShareGroupCommandOptions extends
CommandDefaultOptions {
final OptionSpec<Void> resetOffsetsOpt;
final OptionSpec<Void> deleteOffsetsOpt;
final OptionSpec<Void> dryRunOpt;
+ final OptionSpec<Void> exportOpt;
+ final OptionSpec<Long> resetToOffsetOpt;
+ final OptionSpec<String> resetFromFileOpt;
final OptionSpec<Void> executeOpt;
final OptionSpec<String> resetToDatetimeOpt;
final OptionSpec<Void> resetToEarliestOpt;
final OptionSpec<Void> resetToLatestOpt;
+ final OptionSpec<Void> resetToCurrentOpt;
final OptionSpec<Void> membersOpt;
final OptionSpec<Void> offsetsOpt;
final OptionSpec<String> stateOpt;
@@ -91,7 +100,8 @@ public class ShareGroupCommandOptions extends
CommandDefaultOptions {
final Set<OptionSpec<?>> allGroupSelectionScopeOpts;
final Set<OptionSpec<?>> allTopicSelectionScopeOpts;
final Set<OptionSpec<?>> allShareGroupLevelOpts;
- final Set<OptionSpec<?>> allResetOffsetScenarioOpts;
+ final Set<OptionSpec<?>> allResetOffsetsScopeOpts;
+ final Set<OptionSpec<?>> allResetOffsetsScenarioOpts;
final Set<OptionSpec<?>> allDeleteOffsetsOpts;
public ShareGroupCommandOptions(String[] args) {
@@ -127,12 +137,22 @@ public class ShareGroupCommandOptions extends
CommandDefaultOptions {
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);
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);
membersOpt = parser.accepts("members", MEMBERS_DOC)
.availableIf(describeOpt);
offsetsOpt = parser.accepts("offsets", OFFSETS_DOC)
@@ -147,7 +167,9 @@ public class ShareGroupCommandOptions extends
CommandDefaultOptions {
allGroupSelectionScopeOpts = Set.of(groupOpt, allGroupsOpt);
allTopicSelectionScopeOpts = Set.of(topicOpt, allTopicsOpt);
allShareGroupLevelOpts = Set.of(listOpt, describeOpt, deleteOpt,
resetOffsetsOpt);
- allResetOffsetScenarioOpts = Set.of(resetToDatetimeOpt,
resetToEarliestOpt, resetToLatestOpt);
+ allResetOffsetsScopeOpts = Set.of(topicOpt, allTopicsOpt,
resetFromFileOpt);
+ allResetOffsetsScenarioOpts = Set.of(resetToOffsetOpt,
resetToDatetimeOpt,
+ resetToEarliestOpt, resetToLatestOpt, resetToCurrentOpt,
resetFromFileOpt);
allDeleteOffsetsOpts = Set.of(groupOpt, topicOpt);
options = parser.parse(args);
@@ -176,13 +198,13 @@ public class ShareGroupCommandOptions extends
CommandDefaultOptions {
if (options.has(deleteOpt)) {
if (!options.has(groupOpt) && !options.has(allGroupsOpt))
CommandLineUtils.printUsageAndExit(parser,
- String.format("Option %s takes the options %s or %s",
deleteOpt, groupOpt, allGroupsOpt));
+ String.format("Option %s takes the options %s or %s.",
deleteOpt, groupOpt, allGroupsOpt));
if (options.has(allGroupsOpt) && options.has(groupOpt))
CommandLineUtils.printUsageAndExit(parser,
String.format("Option %s takes either %s or %s, not
both.", deleteOpt, groupOpt, allGroupsOpt));
- if (options.has(topicOpt))
+ if (options.has(allTopicsOpt) || options.has(topicOpt))
CommandLineUtils.printUsageAndExit(parser,
- "Option " + deleteOpt + " does not take the option: " +
topicOpt);
+ String.format("Option %s does not take the options %s or
%s.", deleteOpt, topicOpt, allTopicsOpt));
}
if (options.has(deleteOffsetsOpt)) {
@@ -203,23 +225,30 @@ public class ShareGroupCommandOptions extends
CommandDefaultOptions {
CommandLineUtils.printUsageAndExit(parser,
"Option " + resetOffsetsOpt + " takes the option: " +
groupOpt);
- if (!options.has(topicOpt) && !options.has(allTopicsOpt)) {
+ if (!options.has(topicOpt) && !options.has(allTopicsOpt) &&
!options.has(resetFromFileOpt)) {
CommandLineUtils.printUsageAndExit(parser,
- "Option " + resetOffsetsOpt + " takes one of these
options: " +
allTopicSelectionScopeOpts.stream().map(Object::toString).sorted().collect(Collectors.joining(",
")));
+ "Option " + resetOffsetsOpt + " takes one of these
options: " +
allResetOffsetsScopeOpts.stream().map(Object::toString).sorted().collect(Collectors.joining(",
")));
}
-
- if (!options.has(resetToEarliestOpt) &&
!options.has(resetToLatestOpt) && !options.has(resetToDatetimeOpt)) {
+
+ CommandLineUtils.checkInvalidArgs(parser, options, topicOpt,
minus(allResetOffsetsScopeOpts, topicOpt));
+ CommandLineUtils.checkInvalidArgs(parser, options, allTopicsOpt,
minus(allResetOffsetsScopeOpts, allTopicsOpt));
+ CommandLineUtils.checkInvalidArgs(parser, options,
resetFromFileOpt, minus(allResetOffsetsScopeOpts, resetFromFileOpt));
+
+ if (!options.has(resetToOffsetOpt) &&
!options.has(resetToEarliestOpt) && !options.has(resetToLatestOpt) &&
!options.has(resetToDatetimeOpt) && !options.has(resetToCurrentOpt) &&
!options.has(resetFromFileOpt)) {
CommandLineUtils.printUsageAndExit(parser,
- "Option " + resetOffsetsOpt + " takes one of these
options: " +
allResetOffsetScenarioOpts.stream().map(Object::toString).sorted().collect(Collectors.joining(",
")));
+ "Option " + resetOffsetsOpt + " takes one of these
options: " +
allResetOffsetsScenarioOpts.stream().map(Object::toString).sorted().collect(Collectors.joining(",
")));
}
- CommandLineUtils.checkInvalidArgs(parser, options,
resetToDatetimeOpt, minus(allResetOffsetScenarioOpts, resetToDatetimeOpt));
- CommandLineUtils.checkInvalidArgs(parser, options,
resetToEarliestOpt, minus(allResetOffsetScenarioOpts, resetToEarliestOpt));
- CommandLineUtils.checkInvalidArgs(parser, options,
resetToLatestOpt, minus(allResetOffsetScenarioOpts, resetToLatestOpt));
+ CommandLineUtils.checkInvalidArgs(parser, options,
resetToOffsetOpt, minus(allResetOffsetsScenarioOpts, resetToOffsetOpt));
+ CommandLineUtils.checkInvalidArgs(parser, options,
resetToDatetimeOpt, minus(allResetOffsetsScenarioOpts, resetToDatetimeOpt));
+ CommandLineUtils.checkInvalidArgs(parser, options,
resetToEarliestOpt, minus(allResetOffsetsScenarioOpts, resetToEarliestOpt));
+ CommandLineUtils.checkInvalidArgs(parser, options,
resetToLatestOpt, minus(allResetOffsetsScenarioOpts, resetToLatestOpt));
+ CommandLineUtils.checkInvalidArgs(parser, options,
resetToCurrentOpt, minus(allResetOffsetsScenarioOpts, resetToCurrentOpt));
+ CommandLineUtils.checkInvalidArgs(parser, options,
resetFromFileOpt, minus(allResetOffsetsScenarioOpts, resetFromFileOpt));
}
CommandLineUtils.checkInvalidArgs(parser, options, groupOpt,
minus(allGroupSelectionScopeOpts, groupOpt));
- CommandLineUtils.checkInvalidArgs(parser, options, groupOpt,
minus(allShareGroupLevelOpts, describeOpt, deleteOpt, resetOffsetsOpt));
- CommandLineUtils.checkInvalidArgs(parser, options, topicOpt,
minus(allShareGroupLevelOpts, deleteOpt, resetOffsetsOpt));
+ CommandLineUtils.checkInvalidArgs(parser, options, groupOpt,
minus(allShareGroupLevelOpts, deleteOpt, deleteOffsetsOpt, describeOpt,
resetOffsetsOpt));
+ CommandLineUtils.checkInvalidArgs(parser, options, topicOpt,
minus(allShareGroupLevelOpts, deleteOffsetsOpt, resetOffsetsOpt));
}
}
diff --git
a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java
b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java
index 8a078715fe1..136973ee86f 100644
---
a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java
+++
b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ShareGroupCommandTest.java
@@ -64,6 +64,10 @@ import org.junit.jupiter.api.Test;
import org.mockito.ArgumentMatcher;
import org.mockito.ArgumentMatchers;
+import java.io.BufferedWriter;
+import java.io.File;
+import java.io.FileWriter;
+import java.io.IOException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
@@ -1112,7 +1116,7 @@ public class ShareGroupCommandTest {
}
@Test
- public void testAlterShareGroupMultipleTopicsSuccess() {
+ public void testResetOffsetsMultipleTopicsSuccess() {
String group = "share-group";
String topic1 = "topic1";
String topic2 = "topic2";
@@ -1170,7 +1174,7 @@ public class ShareGroupCommandTest {
}
@Test
- public void testAlterShareGroupToLatestSuccess() {
+ public void testResetOffsetsToLatestSuccess() {
String group = "share-group";
String topic = "topic";
String bootstrapServer = "localhost:9092";
@@ -1223,7 +1227,7 @@ public class ShareGroupCommandTest {
}
@Test
- public void testAlterShareGroupAllTopicsToDatetimeSuccess() {
+ public void testResetOffsetsAllTopicsToDatetimeSuccess() {
String group = "share-group";
String topic1 = "topic1";
String topic2 = "topic2";
@@ -1294,6 +1298,59 @@ public class ShareGroupCommandTest {
}
}
+ @Test
+ public void testResetOffsetsToOffsetSuccess() {
+ String group = "share-group";
+ String topic = "topic";
+ String bootstrapServer = "localhost:9092";
+ String[] cgcArgs = new String[]{"--bootstrap-server", bootstrapServer,
"--reset-offsets", "--to-offset", "10", "--execute", "--topic", topic,
"--group", group};
+ Admin adminClient = mock(KafkaAdminClient.class);
+ TopicPartition t1 = new TopicPartition(topic, 0);
+ TopicPartition t2 = new TopicPartition(topic, 1);
+ ListShareGroupOffsetsResult listShareGroupOffsetsResult =
AdminClientTestUtils.createListShareGroupOffsetsResult(
+ Map.of(
+ group,
+ KafkaFuture.completedFuture(Map.of(
+ t1, new SharePartitionOffsetInfo(10L, Optional.empty(),
Optional.empty()),
+ t2, new SharePartitionOffsetInfo(10L, Optional.empty(),
Optional.empty())))
+ )
+ );
+ Map<String, TopicDescription> descriptions = Map.of(
+ topic, new TopicDescription(topic, false, List.of(
+ new TopicPartitionInfo(0, new Node(0, "localhost", 9092),
List.of(), List.of()),
+ new TopicPartitionInfo(1, new Node(0, "localhost", 9092),
List.of(), List.of()))
+ ));
+ DescribeTopicsResult describeTopicResult =
mock(DescribeTopicsResult.class);
+
when(describeTopicResult.allTopicNames()).thenReturn(completedFuture(descriptions));
+
when(adminClient.describeTopics(anyCollection())).thenReturn(describeTopicResult);
+ when(adminClient.describeTopics(anyCollection(),
any(DescribeTopicsOptions.class))).thenReturn(describeTopicResult);
+ when(adminClient.listShareGroupOffsets(any(),
any(ListShareGroupOffsetsOptions.class))).thenReturn(listShareGroupOffsetsResult);
+
+ AlterShareGroupOffsetsResult alterShareGroupOffsetsResult =
mockAlterShareGroupOffsets(adminClient, group);
+ Map<TopicPartition, OffsetAndMetadata> partitionOffsets = Map.of(t1,
new OffsetAndMetadata(40L), t2, new OffsetAndMetadata(40L));
+ ListOffsetsResult listOffsetsResult =
AdminClientTestUtils.createListOffsetsResult(partitionOffsets);
+ when(adminClient.listOffsets(any(),
any(ListOffsetsOptions.class))).thenReturn(listOffsetsResult);
+
+ ShareGroupDescription exp = new ShareGroupDescription(
+ group,
+ List.of(),
+ GroupState.EMPTY,
+ new Node(0, "host1", 9090), 0, 0);
+ DescribeShareGroupsResult describeShareGroupsResult =
mock(DescribeShareGroupsResult.class);
+
when(describeShareGroupsResult.describedGroups()).thenReturn(Map.of(group,
KafkaFuture.completedFuture(exp)));
+ when(adminClient.describeShareGroups(ArgumentMatchers.anyCollection(),
any(DescribeShareGroupsOptions.class))).thenReturn(describeShareGroupsResult);
+ Function<Collection<TopicPartition>,
ArgumentMatcher<Map<TopicPartition, OffsetSpec>>> offsetsArgMatcher =
expectedPartitions ->
+ topicPartitionOffsets -> topicPartitionOffsets != null &&
topicPartitionOffsets.keySet().equals(expectedPartitions) &&
+ topicPartitionOffsets.values().stream().allMatch(offsetSpec ->
offsetSpec instanceof OffsetSpec.LatestSpec);
+ try (ShareGroupService service = getShareGroupService(cgcArgs,
adminClient)) {
+ service.resetOffsets();
+ verify(adminClient).alterShareGroupOffsets(eq(group), anyMap(),
any(AlterShareGroupOffsetsOptions.class));
+ verify(adminClient,
times(1)).listOffsets(ArgumentMatchers.argThat(offsetsArgMatcher.apply(Set.of(t1,
t2))), any());
+ verify(alterShareGroupOffsetsResult, times(1)).all();
+
verify(adminClient).describeShareGroups(ArgumentMatchers.anyCollection(),
any(DescribeShareGroupsOptions.class));
+ }
+ }
+
@Test
public void testResetOffsetsDryRunSuccess() {
String group = "share-group";
@@ -1342,7 +1399,7 @@ public class ShareGroupCommandTest {
}
@Test
- public void testAlterShareGroupOffsetsFailureWithoutTopic() {
+ public void testResetOffsetsFailureWithoutTopic() {
String bootstrapServer = "localhost:9092";
String group = "share-group";
Admin adminClient = mock(KafkaAdminClient.class);
@@ -1351,7 +1408,7 @@ public class ShareGroupCommandTest {
AtomicBoolean exited = new AtomicBoolean(false);
Exit.setExitProcedure(((statusCode, message) -> {
assertNotEquals(0, statusCode);
- assertTrue(message.contains("Option [reset-offsets] takes one of
these options: [all-topics], [topic]"));
+ assertTrue(message.contains("Option [reset-offsets] takes one of
these options: [all-topics], [from-file], [topic]"));
exited.set(true);
}));
try {
@@ -1362,7 +1419,7 @@ public class ShareGroupCommandTest {
}
@Test
- public void testAlterShareGroupOffsetsFailureWithNonEmptyGroup() {
+ public void testResetOffsetsFailureWithNonEmptyGroup() {
String group = "share-group";
String topic = "topic";
String bootstrapServer = "localhost:9092";
@@ -1401,7 +1458,7 @@ public class ShareGroupCommandTest {
}
@Test
- public void testAlterShareGroupOffsetsArgsFailureWithoutResetOffsetsArgs()
{
+ public void testResetOffsetsArgsFailureWithoutResetOffsetsArgs() {
String bootstrapServer = "localhost:9092";
String group = "share-group";
Admin adminClient = mock(KafkaAdminClient.class);
@@ -1410,7 +1467,7 @@ public class ShareGroupCommandTest {
AtomicBoolean exited = new AtomicBoolean(false);
Exit.setExitProcedure(((statusCode, message) -> {
assertNotEquals(0, statusCode);
- assertTrue(message.contains("Option [reset-offsets] takes one of
these options: [to-datetime], [to-earliest], [to-latest]"));
+ assertTrue(message.contains("Option [reset-offsets] takes one of
these options: [from-file], [to-current], [to-datetime], [to-earliest],
[to-latest], [to-offset]"));
exited.set(true);
}));
try {
@@ -1421,7 +1478,7 @@ public class ShareGroupCommandTest {
}
@Test
- public void testAlterShareGroupUnsubscribedTopicSuccess() {
+ public void testResetOffsetsUnsubscribedTopicSuccess() {
String group = "share-group";
String topic = "none";
String bootstrapServer = "localhost:9092";
@@ -1468,7 +1525,7 @@ public class ShareGroupCommandTest {
}
@Test
- public void testAlterShareGroupNonExistentGroupSuccess() {
+ public void testResetOffsetsNonExistentGroupSuccess() {
String group = "share-group";
String topic = "none";
String bootstrapServer = "localhost:9092";
@@ -1511,6 +1568,136 @@ public class ShareGroupCommandTest {
}
}
+ @Test
+ public void testResetOffsetsExportImportPlan() throws Exception {
+ String group = "share-group";
+ String topic = "topic";
+ String bootstrapServer = "localhost:9092";
+ String[] cgcArgs = new String[]{"--bootstrap-server", bootstrapServer,
"--reset-offsets", "--to-current", "--dry-run", "--topic", topic, "--group",
group};
+ File file = TestUtils.tempFile("reset", ".csv");
+ Admin adminClient = mock(KafkaAdminClient.class);
+
+ ListShareGroupOffsetsResult listShareGroupOffsetsResult =
AdminClientTestUtils.createListShareGroupOffsetsResult(
+ Map.of(
+ group,
+ KafkaFuture.completedFuture(Map.of(new TopicPartition("topic",
0), new SharePartitionOffsetInfo(10L, Optional.empty(), Optional.empty())))
+ )
+ );
+ when(adminClient.listShareGroupOffsets(any(),
any(ListShareGroupOffsetsOptions.class))).thenReturn(listShareGroupOffsetsResult);
+ // Because there is a single topic-partition and there is a known
start offset, the first ListOffsetsResults is empty.
+ Map<TopicPartition, OffsetAndMetadata> partitionEmptyOffsets =
Map.of();
+ ListOffsetsResult listEmptyOffsetsResult =
AdminClientTestUtils.createListOffsetsResult(partitionEmptyOffsets);
+ Map<TopicPartition, OffsetAndMetadata> partitionStartOffsets = Map.of(
+ new TopicPartition(topic, 0), new OffsetAndMetadata(5));
+ ListOffsetsResult listStartOffsetsResult =
AdminClientTestUtils.createListOffsetsResult(partitionStartOffsets);
+ Map<TopicPartition, OffsetAndMetadata> partitionEndOffsets = Map.of(
+ new TopicPartition(topic, 0), new OffsetAndMetadata(25));
+ ListOffsetsResult listEndOffsetsResult =
AdminClientTestUtils.createListOffsetsResult(partitionEndOffsets);
+ when(adminClient.listOffsets(any(),
any(ListOffsetsOptions.class))).thenReturn(listEmptyOffsetsResult).thenReturn(listStartOffsetsResult).thenReturn(listEndOffsetsResult);
+
+ AlterShareGroupOffsetsResult alterShareGroupOffsetsResult =
mockAlterShareGroupOffsets(adminClient, group);
+
+ ShareGroupDescription exp = new ShareGroupDescription(
+ group,
+ List.of(),
+ GroupState.EMPTY,
+ new Node(0, "host1", 9090), 0, 0);
+ DescribeShareGroupsResult describeShareGroupsResult =
mock(DescribeShareGroupsResult.class);
+
when(describeShareGroupsResult.describedGroups()).thenReturn(Map.of(group,
KafkaFuture.completedFuture(exp)));
+ when(adminClient.describeShareGroups(any(),
any(DescribeShareGroupsOptions.class))).thenReturn(describeShareGroupsResult);
+ Map<String, TopicDescription> descriptions = Map.of(
+ topic, new TopicDescription(topic, false, List.of(
+ new TopicPartitionInfo(0, new Node(0, "localhost", 9092),
List.of(), List.of())
+ )));
+ DescribeTopicsResult describeTopicResult =
mock(DescribeTopicsResult.class);
+
when(describeTopicResult.allTopicNames()).thenReturn(completedFuture(descriptions));
+
when(adminClient.describeTopics(anyCollection())).thenReturn(describeTopicResult);
+ when(adminClient.describeTopics(anyCollection(),
any(DescribeTopicsOptions.class))).thenReturn(describeTopicResult);
+
+ try (ShareGroupService service = getShareGroupService(cgcArgs,
adminClient)) {
+ Map<String, Map<TopicPartition, OffsetAndMetadata>>
offsetsForGroups = service.resetOffsets();
+ writeContentToFile(file,
service.exportOffsetsToCsv(offsetsForGroups));
+
+ String[] cgcArgsExec = new String[]{"--bootstrap-server",
bootstrapServer, "--reset-offsets", "--from-file", file.getCanonicalPath(),
"--execute", "--topic", topic, "--group", group};
+ try (ShareGroupService serviceExec =
getShareGroupService(cgcArgsExec, adminClient)) {
+ Map<String, Map<TopicPartition, OffsetAndMetadata>>
offsetsForGroupsExec = serviceExec.resetOffsets();
+ assertEquals(1, offsetsForGroupsExec.size());
+ assertEquals(1, offsetsForGroupsExec.get(group).size());
+ assertEquals(10, offsetsForGroupsExec.get(group).get(new
TopicPartition(topic, 0)).offset());
+ verify(adminClient).alterShareGroupOffsets(eq(group),
anyMap(), any(AlterShareGroupOffsetsOptions.class));
+ verify(adminClient, times(6)).describeTopics(anyCollection(),
any(DescribeTopicsOptions.class));
+ verify(alterShareGroupOffsetsResult, times(1)).all();
+ verify(adminClient,
times(2)).describeShareGroups(ArgumentMatchers.anyCollection(),
any(DescribeShareGroupsOptions.class));
+ }
+ }
+ }
+
+ @Test
+ public void testResetOffsetsExportImportComplexPlan() throws Exception {
+ String group = "share-group";
+ String topic = "topic";
+ String topic2 = "topic2";
+ String bootstrapServer = "localhost:9092";
+ File file = TestUtils.tempFile("reset", ".csv");
+ Admin adminClient = mock(KafkaAdminClient.class);
+
+ // This test uses the GROUP,TOPIC,PARTITION,OFFSET format which is
generated by the kafka-consumer-groups.sh tool when
+ // multiple groups are being exported. The reset is filtered by the
group ID.
+ FileWriter writer = new FileWriter(file);
+ writer.write("share-group,topic,0,10\n");
+ writer.write("share-group,topic2,0,20\n");
+ writer.write("share-group2,topic,0,30\n");
+ writer.close();
+
+ ListShareGroupOffsetsResult listShareGroupOffsetsResult =
AdminClientTestUtils.createListShareGroupOffsetsResult(
+ Map.of(
+ group,
+ KafkaFuture.completedFuture(Map.of(new TopicPartition("topic",
0), new SharePartitionOffsetInfo(0L, Optional.empty(), Optional.empty())))
+ )
+ );
+ when(adminClient.listShareGroupOffsets(any(),
any(ListShareGroupOffsetsOptions.class))).thenReturn(listShareGroupOffsetsResult);
+ Map<TopicPartition, OffsetAndMetadata> partitionStartOffsets = Map.of(
+ new TopicPartition(topic, 0), new OffsetAndMetadata(5),
+ new TopicPartition(topic2, 0), new OffsetAndMetadata(5));
+ ListOffsetsResult listStartOffsetsResult =
AdminClientTestUtils.createListOffsetsResult(partitionStartOffsets);
+ Map<TopicPartition, OffsetAndMetadata> partitionEndOffsets = Map.of(
+ new TopicPartition(topic, 0), new OffsetAndMetadata(25),
+ new TopicPartition(topic2, 0), new OffsetAndMetadata(25));
+ ListOffsetsResult listEndOffsetsResult =
AdminClientTestUtils.createListOffsetsResult(partitionEndOffsets);
+ when(adminClient.listOffsets(any(),
any(ListOffsetsOptions.class))).thenReturn(listStartOffsetsResult).thenReturn(listEndOffsetsResult);
+
+ AlterShareGroupOffsetsResult alterShareGroupOffsetsResult =
mockAlterShareGroupOffsets(adminClient, group);
+
+ ShareGroupDescription exp = new ShareGroupDescription(
+ group,
+ List.of(),
+ GroupState.EMPTY,
+ new Node(0, "host1", 9090), 0, 0);
+ DescribeShareGroupsResult describeShareGroupsResult =
mock(DescribeShareGroupsResult.class);
+
when(describeShareGroupsResult.describedGroups()).thenReturn(Map.of(group,
KafkaFuture.completedFuture(exp)));
+ when(adminClient.describeShareGroups(any(),
any(DescribeShareGroupsOptions.class))).thenReturn(describeShareGroupsResult);
+ Map<String, TopicDescription> descriptions = Map.of(
+ topic, new TopicDescription(topic, false, List.of(
+ new TopicPartitionInfo(0, new Node(0, "localhost", 9092),
List.of(), List.of())
+ )));
+ DescribeTopicsResult describeTopicResult =
mock(DescribeTopicsResult.class);
+
when(describeTopicResult.allTopicNames()).thenReturn(completedFuture(descriptions));
+
when(adminClient.describeTopics(anyCollection())).thenReturn(describeTopicResult);
+ when(adminClient.describeTopics(anyCollection(),
any(DescribeTopicsOptions.class))).thenReturn(describeTopicResult);
+
+ String[] cgcArgsExec = new String[]{"--bootstrap-server",
bootstrapServer, "--reset-offsets", "--from-file", file.getCanonicalPath(),
"--execute", "--all-topics", "--group", group};
+ try (ShareGroupService serviceExec = getShareGroupService(cgcArgsExec,
adminClient)) {
+ Map<String, Map<TopicPartition, OffsetAndMetadata>>
offsetsForGroupsExec = serviceExec.resetOffsets();
+ assertEquals(1, offsetsForGroupsExec.size());
+ assertEquals(2, offsetsForGroupsExec.get(group).size());
+ assertEquals(10, offsetsForGroupsExec.get(group).get(new
TopicPartition(topic, 0)).offset());
+ verify(adminClient).alterShareGroupOffsets(eq(group), anyMap(),
any(AlterShareGroupOffsetsOptions.class));
+ verify(adminClient, times(2)).describeTopics(anyCollection(),
any(DescribeTopicsOptions.class));
+ verify(alterShareGroupOffsetsResult, times(1)).all();
+ verify(adminClient,
times(1)).describeShareGroups(ArgumentMatchers.anyCollection(),
any(DescribeShareGroupsOptions.class));
+ }
+ }
+
private AlterShareGroupOffsetsResult mockAlterShareGroupOffsets(Admin
client, String groupId) {
AlterShareGroupOffsetsResult alterShareGroupOffsetsResult =
mock(AlterShareGroupOffsetsResult.class);
KafkaFutureImpl<Void> resultFuture = new KafkaFutureImpl<>();
@@ -1546,6 +1733,12 @@ public class ShareGroupCommandTest {
return () -> Assertions.assertDoesNotThrow(service::deleteOffsets);
}
+ private void writeContentToFile(File file, String content) throws
IOException {
+ try (BufferedWriter bw = new BufferedWriter(new FileWriter(file))) {
+ bw.write(content);
+ }
+ }
+
private boolean checkArgsHeaderOutput(List<String> args, String output) {
if (!output.contains("GROUP")) {
return false;