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..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 @@ -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; 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..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 @@ -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; @@ -293,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)); @@ -427,6 +439,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);