Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -881,6 +881,7 @@ private Collection<TopicPartition> getPartitionsToReset(String groupId) throws E
List<String> topics = opts.options.valuesOf(opts.inputTopicOpt);

List<TopicPartition> 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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -65,13 +66,15 @@

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;
import static org.junit.jupiter.api.Assertions.assertNotNull;
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;
Expand Down Expand Up @@ -293,21 +296,30 @@ public void testGroupStatesFromString() {
@Test
public void testAdminRequestsForResetOffsets() {
Admin adminClient = mock(KafkaAdminClient.class);
String topic = "topic1";
String groupId = "foo-group";
List<String> args = List.of("--bootstrap-server", "localhost:9092", "--group", groupId, "--reset-offsets", "--input-topic", "topic1", "--to-latest");
List<String> topics = List.of("topic1");
List<String> topics = List.of(topic);

DescribeTopicsResult describeTopicsResult = mock(DescribeTopicsResult.class);
when(adminClient.describeStreamsGroups(List.of(groupId)))
.thenReturn(describeStreamsResult(groupId, GroupState.DEAD));
Map<String, TopicDescription> 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<TopicPartition, OffsetAndMetadata> 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));

Expand Down Expand Up @@ -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<String> 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<String, TopicDescription> 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<TopicPartition, OffsetAndMetadata> 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);
Expand Down