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;

Reply via email to