From f9a87f14990961fc9f8f2434ee12c9236ee1c803 Mon Sep 17 00:00:00 2001 From: JimmyWang6 Date: Thu, 31 Jul 2025 17:50:27 +0800 Subject: [PATCH 1/4] apply fix #20168 for altering offset for stream group --- .../org/apache/kafka/tools/streams/StreamsGroupCommand.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java b/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java index 0f68bf8290053..85fbf75d3a40e 100644 --- a/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java @@ -893,6 +893,8 @@ private Collection getPartitionsToReset(String groupId) throws E } private Map prepareOffsetsToReset(String groupId, Collection partitionsToReset) { + offsetsUtils.checkAllTopicPartitionsValid(partitionsToReset); + if (opts.options.has(opts.resetToOffsetOpt)) { return offsetsUtils.resetToOffset(partitionsToReset); } else if (opts.options.has(opts.resetToEarliestOpt)) { From 59a5c60da926f070a0d34023b7c29b0f25474e92 Mon Sep 17 00:00:00 2001 From: JimmyWang6 Date: Mon, 18 Aug 2025 14:25:23 +0800 Subject: [PATCH 2/4] add tests --- .../tools/streams/StreamsGroupCommand.java | 2 +- .../consumer/group/ShareGroupCommandTest.java | 2 +- .../streams/StreamsGroupCommandTest.java | 40 +++++++++++++++++++ 3 files changed, 42 insertions(+), 2 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java b/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java index 85fbf75d3a40e..33859eb88abab 100644 --- a/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java @@ -881,6 +881,7 @@ private Collection getPartitionsToReset(String groupId) throws E List topics = opts.options.valuesOf(opts.inputTopicOpt); List partitions = offsetsUtils.parseTopicPartitionsToReset(topics); + offsetsUtils.checkAllTopicPartitionsValid(partitions); // if the user specified topics that do not belong to this group, we filter them out partitions = filterExistingGroupTopics(groupId, partitions); return partitions; @@ -893,7 +894,6 @@ private Collection getPartitionsToReset(String groupId) throws E } private Map prepareOffsetsToReset(String groupId, Collection partitionsToReset) { - offsetsUtils.checkAllTopicPartitionsValid(partitionsToReset); if (opts.options.has(opts.resetToOffsetOpt)) { return offsetsUtils.resetToOffset(partitionsToReset); 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 b14d66c652ab6..9333bbbb65e12 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 @@ -1373,7 +1373,7 @@ public void testAlterShareGroupOffsetsArgsFailureWithoutResetOffsetsArgs() { } @Test - public void testAlterShareGroupFailureFailureWithNonExistentTopic() { + public void testAlterShareGroupFailureWithNonExistentTopic() { String group = "share-group"; String topic = "none"; String bootstrapServer = "localhost:9092"; diff --git a/tools/src/test/java/org/apache/kafka/tools/streams/StreamsGroupCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/streams/StreamsGroupCommandTest.java index 6f38c47f15aa2..77f41abf425d5 100644 --- a/tools/src/test/java/org/apache/kafka/tools/streams/StreamsGroupCommandTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/streams/StreamsGroupCommandTest.java @@ -43,6 +43,7 @@ import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.TopicPartitionInfo; +import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; import org.apache.kafka.common.internals.KafkaFutureImpl; import org.apache.kafka.test.TestUtils; @@ -65,6 +66,7 @@ import joptsimple.OptionException; +import static org.apache.kafka.common.KafkaFuture.completedFuture; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertInstanceOf; @@ -72,6 +74,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyCollection; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.times; @@ -427,6 +430,43 @@ public void testDeleteNonStreamsGroup() { service.close(); } + + @Test + public void testResetOffsetsWithPartitionNotExist() { + Admin adminClient = mock(KafkaAdminClient.class); + String groupId = "foo-group"; + String topic = "topic"; + List args = new ArrayList<>(Arrays.asList("--bootstrap-server", "localhost:9092", "--group", groupId, "--reset-offsets", "--input-topic", "topic:3", "--to-latest")); + + when(adminClient.describeStreamsGroups(List.of(groupId))) + .thenReturn(describeStreamsResult(groupId, GroupState.DEAD)); + DescribeTopicsResult describeTopicsResult = mock(DescribeTopicsResult.class); + + Map descriptions = Map.of( + topic, new TopicDescription(topic, false, List.of( + new TopicPartitionInfo(0, Node.noNode(), List.of(), List.of())) + )); + when(adminClient.describeTopics(anyCollection())) + .thenReturn(describeTopicsResult); + when(adminClient.describeTopics(eq(List.of(topic)), any(DescribeTopicsOptions.class))) + .thenReturn(describeTopicsResult); + when(describeTopicsResult.allTopicNames()).thenReturn(completedFuture(descriptions)); + when(adminClient.listOffsets(any(), any())) + .thenReturn(listOffsetsResult()); + ListStreamsGroupOffsetsResult result = mock(ListStreamsGroupOffsetsResult.class); + Map committedOffsetsMap = Map.of( + new TopicPartition(topic, 0), + new OffsetAndMetadata(12, Optional.of(0), ""), + new TopicPartition(topic, 1), + new OffsetAndMetadata(12, Optional.of(0), "") + ); + + when(adminClient.listStreamsGroupOffsets(ArgumentMatchers.anyMap())).thenReturn(result); + when(result.partitionsToOffsetAndMetadata(ArgumentMatchers.anyString())).thenReturn(KafkaFuture.completedFuture(committedOffsetsMap)); + StreamsGroupCommand.StreamsGroupService service = getStreamsGroupService(args.toArray(new String[0]), adminClient); + assertThrows(UnknownTopicOrPartitionException.class, () -> service.resetOffsets()); + service.close(); + } private ListGroupsResult listGroupResult(String groupId) { ListGroupsResult listGroupsResult = mock(ListGroupsResult.class); From 111d69e904a898ddb5f3f5acf3fa5c8f9d47794c Mon Sep 17 00:00:00 2001 From: JimmyWang6 Date: Mon, 18 Aug 2025 14:27:30 +0800 Subject: [PATCH 3/4] remove unrelated change --- .../java/org/apache/kafka/tools/streams/StreamsGroupCommand.java | 1 - 1 file changed, 1 deletion(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java b/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java index 33859eb88abab..0c54f6c53f99d 100644 --- a/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/streams/StreamsGroupCommand.java @@ -894,7 +894,6 @@ private Collection getPartitionsToReset(String groupId) throws E } private Map prepareOffsetsToReset(String groupId, Collection partitionsToReset) { - if (opts.options.has(opts.resetToOffsetOpt)) { return offsetsUtils.resetToOffset(partitionsToReset); } else if (opts.options.has(opts.resetToEarliestOpt)) { From 30eb73c839c4411841af27196215311eab158131 Mon Sep 17 00:00:00 2001 From: JimmyWang6 Date: Mon, 25 Aug 2025 18:23:17 +0800 Subject: [PATCH 4/4] fix test --- .../tools/streams/StreamsGroupCommandTest.java | 15 ++++++++++++--- 1 file changed, 12 insertions(+), 3 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/streams/StreamsGroupCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/streams/StreamsGroupCommandTest.java index 77f41abf425d5..4f1e116437e63 100644 --- a/tools/src/test/java/org/apache/kafka/tools/streams/StreamsGroupCommandTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/streams/StreamsGroupCommandTest.java @@ -296,21 +296,30 @@ public void testGroupStatesFromString() { @Test public void testAdminRequestsForResetOffsets() { Admin adminClient = mock(KafkaAdminClient.class); + String topic = "topic1"; String groupId = "foo-group"; List args = List.of("--bootstrap-server", "localhost:9092", "--group", groupId, "--reset-offsets", "--input-topic", "topic1", "--to-latest"); - List topics = List.of("topic1"); + List topics = List.of(topic); + DescribeTopicsResult describeTopicsResult = mock(DescribeTopicsResult.class); when(adminClient.describeStreamsGroups(List.of(groupId))) .thenReturn(describeStreamsResult(groupId, GroupState.DEAD)); + Map descriptions = Map.of( + topic, new TopicDescription(topic, false, List.of( + new TopicPartitionInfo(0, Node.noNode(), List.of(), List.of())) + )); + when(adminClient.describeTopics(anyCollection())) + .thenReturn(describeTopicsResult); when(adminClient.describeTopics(eq(topics), any(DescribeTopicsOptions.class))) - .thenReturn(describeTopicsResult(topics, 1)); + .thenReturn(describeTopicsResult); + when(describeTopicsResult.allTopicNames()).thenReturn(completedFuture(descriptions)); when(adminClient.listOffsets(any(), any())) .thenReturn(listOffsetsResult()); ListGroupsResult listGroupsResult = listGroupResult(groupId); when(adminClient.listGroups(any(ListGroupsOptions.class))).thenReturn(listGroupsResult); ListStreamsGroupOffsetsResult result = mock(ListStreamsGroupOffsetsResult.class); Map committedOffsetsMap = new HashMap<>(); - committedOffsetsMap.put(new TopicPartition("topic1", 0), mock(OffsetAndMetadata.class)); + committedOffsetsMap.put(new TopicPartition(topic, 0), mock(OffsetAndMetadata.class)); when(adminClient.listStreamsGroupOffsets(ArgumentMatchers.anyMap())).thenReturn(result); when(result.partitionsToOffsetAndMetadata(ArgumentMatchers.anyString())).thenReturn(KafkaFuture.completedFuture(committedOffsetsMap));