From 88e7a0ca4f15c3816c9f49ef1e4ffc0c1d36f4cc Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Fri, 1 Nov 2019 16:28:08 +0000 Subject: [PATCH 01/14] KAFKA-9130: KIP-518 Allow listing consumer groups per state Co-authored-by: Mickael Maison Co-authored-by: Edoardo Comar --- .../clients/admin/ConsumerGroupListing.java | 40 ++++++++++- .../kafka/clients/admin/KafkaAdminClient.java | 8 ++- .../admin/ListConsumerGroupsOptions.java | 35 ++++++++++ .../common/requests/ListGroupsRequest.java | 4 ++ .../common/message/ListGroupsRequest.json | 8 ++- .../common/message/ListGroupsResponse.json | 8 ++- .../clients/admin/KafkaAdminClientTest.java | 37 ++++++++++ .../kafka/admin/ConsumerGroupCommand.scala | 67 ++++++++++++++++--- .../coordinator/group/GroupCoordinator.scala | 5 +- .../coordinator/group/GroupMetadata.scala | 5 +- .../main/scala/kafka/server/KafkaApis.scala | 28 +++++--- ...aslClientsWithInvalidCredentialsTest.scala | 2 +- .../admin/DeleteConsumerGroupsTest.scala | 10 +-- .../admin/DescribeConsumerGroupTest.scala | 8 --- .../kafka/admin/ListConsumerGroupTest.scala | 53 ++++++++++++++- .../group/GroupCoordinatorTest.scala | 56 ++++++++++++++-- 16 files changed, 325 insertions(+), 49 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java b/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java index 46da9628010b6..205e7ebacdcfd 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java @@ -17,22 +17,29 @@ package org.apache.kafka.clients.admin; +import java.util.Objects; + +import org.apache.kafka.common.ConsumerGroupState; + /** * A listing of a consumer group in the cluster. */ public class ConsumerGroupListing { private final String groupId; private final boolean isSimpleConsumerGroup; + private final ConsumerGroupState state; /** * Create an instance with the specified parameters. * * @param groupId Group Id * @param isSimpleConsumerGroup If consumer group is simple or not. + * @param state The state of the consumer group */ - public ConsumerGroupListing(String groupId, boolean isSimpleConsumerGroup) { + public ConsumerGroupListing(String groupId, boolean isSimpleConsumerGroup, ConsumerGroupState state) { this.groupId = groupId; this.isSimpleConsumerGroup = isSimpleConsumerGroup; + this.state = state; } /** @@ -49,11 +56,42 @@ public boolean isSimpleConsumerGroup() { return isSimpleConsumerGroup; } + /** + * Consumer Group state + */ + public ConsumerGroupState state() { + return state; + } + @Override public String toString() { return "(" + "groupId='" + groupId + '\'' + ", isSimpleConsumerGroup=" + isSimpleConsumerGroup + + ", state=" + state + ')'; } + + @Override + public int hashCode() { + return Objects.hash(groupId, isSimpleConsumerGroup, state); + } + + @Override + public boolean equals(final Object o) { + if (this == o) return true; + if (o == null || getClass() != o.getClass()) return false; + final ConsumerGroupListing that = (ConsumerGroupListing) o; + if (groupId == null) { + if (that.groupId != null) + return false; + } else if (!groupId.equals(that.groupId)) + return false; + if (isSimpleConsumerGroup != that.isSimpleConsumerGroup) + return false; + if (state != that.state) + return false; + return true; + } + } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index e0e6cb6cffa00..0e37acad5770c 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3052,14 +3052,18 @@ void handleResponse(AbstractResponse abstractResponse) { runnable.call(new Call("listConsumerGroups", deadline, new ConstantNodeIdProvider(node.id())) { @Override ListGroupsRequest.Builder createRequest(int timeoutMs) { - return new ListGroupsRequest.Builder(new ListGroupsRequestData()); + List states = options.states().orElse(Collections.emptySet()) + .stream() + .map(s -> s.toString()) + .collect(Collectors.toList()); + return new ListGroupsRequest.Builder(new ListGroupsRequestData().setStates(states)); } private void maybeAddConsumerGroup(ListGroupsResponseData.ListedGroup group) { String protocolType = group.protocolType(); if (protocolType.equals(ConsumerProtocol.PROTOCOL_TYPE) || protocolType.isEmpty()) { final String groupId = group.groupId(); - final ConsumerGroupListing groupListing = new ConsumerGroupListing(groupId, protocolType.isEmpty()); + final ConsumerGroupListing groupListing = new ConsumerGroupListing(groupId, protocolType.isEmpty(), ConsumerGroupState.parse(group.groupState())); results.addListing(groupListing); } } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java index eb27c795f2afe..015a7add496ee 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java @@ -17,6 +17,11 @@ package org.apache.kafka.clients.admin; +import java.util.EnumSet; +import java.util.Optional; +import java.util.Set; + +import org.apache.kafka.common.ConsumerGroupState; import org.apache.kafka.common.annotation.InterfaceStability; /** @@ -26,4 +31,34 @@ */ @InterfaceStability.Evolving public class ListConsumerGroupsOptions extends AbstractOptions { + + private Optional> states = Optional.empty(); + + /** + * Only groups in these states will be returned by listConsumerGroups() + * If not set, all groups are returned without their states + * throw IllegalArgumentException if states is empty + */ + public ListConsumerGroupsOptions inStates(Set states) { + if (states == null || states.isEmpty()) { + throw new IllegalArgumentException("states should not be null or empty"); + } + this.states = Optional.of(states); + return this; + } + + /** + * All groups with their states will be returned by listConsumerGroups() + */ + public ListConsumerGroupsOptions inAnyState() { + this.states = Optional.of(EnumSet.allOf(ConsumerGroupState.class)); + return this; + } + + /** + * Returns the list of States that are requested + */ + public Optional> states() { + return states; + } } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java index caa6f47906bfd..1b48aae6d1f8a 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java @@ -66,6 +66,10 @@ public ListGroupsRequest(Struct struct, short version) { this.data = new ListGroupsRequestData(struct, version); } + public ListGroupsRequestData data() { + return data; + } + @Override public ListGroupsResponse getErrorResponse(int throttleTimeMs, Throwable e) { ListGroupsResponseData listGroupsResponseData = new ListGroupsResponseData(). diff --git a/clients/src/main/resources/common/message/ListGroupsRequest.json b/clients/src/main/resources/common/message/ListGroupsRequest.json index f0130e264902b..513d822c7e803 100644 --- a/clients/src/main/resources/common/message/ListGroupsRequest.json +++ b/clients/src/main/resources/common/message/ListGroupsRequest.json @@ -20,8 +20,14 @@ // Version 1 and 2 are the same as version 0. // // Version 3 is the first flexible version. - "validVersions": "0-3", + // + // Version 4 adds the States flexible field (KIP-518). + "validVersions": "0-4", "flexibleVersions": "3+", "fields": [ + { "name": "States", "type": "[]string", "versions": "4+", "tag": 0, "taggedVersions": "4+", + "about": "The states of the groups we want to list", "fields": [ + { "name": "Name", "type": "string", "versions": "3+", "about": "The group state" } + ]} ] } diff --git a/clients/src/main/resources/common/message/ListGroupsResponse.json b/clients/src/main/resources/common/message/ListGroupsResponse.json index aa8bba6ebc738..eb8914edbc42f 100644 --- a/clients/src/main/resources/common/message/ListGroupsResponse.json +++ b/clients/src/main/resources/common/message/ListGroupsResponse.json @@ -22,7 +22,9 @@ // Starting in version 2, on quota violation, brokers send out responses before throttling. // // Version 3 is the first flexible version. - "validVersions": "0-3", + // + // Version 4 adds the GroupState flexible field (KIP-518). + "validVersions": "0-4", "flexibleVersions": "3+", "fields": [ { "name": "ThrottleTimeMs", "type": "int32", "versions": "1+", "ignorable": true, @@ -34,7 +36,9 @@ { "name": "GroupId", "type": "string", "versions": "0+", "entityType": "groupId", "about": "The group ID." }, { "name": "ProtocolType", "type": "string", "versions": "0+", - "about": "The group protocol type." } + "about": "The group protocol type." }, + { "name": "GroupState", "type": "string", "versions": "4+", "tag": 0, "taggedVersions": "4+", "ignorable": true, + "about": "The group state string." } ]} ] } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index 2510ea727d43d..2dc8274ab5d36 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -26,6 +26,7 @@ import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.clients.consumer.internals.ConsumerProtocol; import org.apache.kafka.common.Cluster; +import org.apache.kafka.common.ConsumerGroupState; import org.apache.kafka.common.ElectionType; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.KafkaFuture; @@ -1282,6 +1283,7 @@ public void testListConsumerGroups() throws Exception { Set groupIds = new HashSet<>(); for (ConsumerGroupListing listing : listings) { groupIds.add(listing.groupId()); + assertEquals(ConsumerGroupState.UNKNOWN, listing.state()); } assertEquals(Utils.mkSet("group-1", "group-2", "group-3"), groupIds); @@ -1312,6 +1314,41 @@ public void testListConsumerGroupsMetadataFailure() throws Exception { } } + @Test + public void testListConsumerGroupsWithStates() throws Exception { + try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(mockCluster(1, 0))) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + + env.kafkaClient().prepareResponse(prepareMetadataResponse(env.cluster(), Errors.NONE)); + + env.kafkaClient().prepareResponseFrom( + new ListGroupsResponse( + new ListGroupsResponseData() + .setErrorCode(Errors.NONE.code()) + .setGroups(Arrays.asList( + new ListGroupsResponseData.ListedGroup() + .setGroupId("group-1") + .setProtocolType(ConsumerProtocol.PROTOCOL_TYPE) + .setGroupState("Stable"), + new ListGroupsResponseData.ListedGroup() + .setGroupId("group-2") + .setGroupState("Empty") + ))), + env.cluster().nodeById(0)); + + final ListConsumerGroupsOptions options = new ListConsumerGroupsOptions().inAnyState(); + final ListConsumerGroupsResult result = env.adminClient().listConsumerGroups(options); + Collection listings = result.valid().get(); + + assertEquals(2, listings.size()); + List expected = new ArrayList<>(); + expected.add(new ConsumerGroupListing("group-2", true, ConsumerGroupState.EMPTY)); + expected.add(new ConsumerGroupListing("group-1", false, ConsumerGroupState.STABLE)); + assertEquals(expected, listings); + assertEquals(0, result.errors().get().size()); + } + } + @Test public void testOffsetCommitNumRetries() throws Exception { final Cluster cluster = mockCluster(3, 0); diff --git a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala index 34d8bbb774609..1e22f30c9f61a 100755 --- a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala +++ b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala @@ -41,9 +41,12 @@ import org.apache.kafka.common.protocol.Errors import scala.collection.immutable.TreeMap import scala.reflect.ClassTag import org.apache.kafka.common.requests.ListOffsetResponse +import org.apache.kafka.common.ConsumerGroupState object ConsumerGroupCommand extends Logging { + val allStates = ConsumerGroupState.values.toList + def main(args: Array[String]): Unit = { val opts = new ConsumerGroupCommandOptions(args) @@ -60,7 +63,7 @@ object ConsumerGroupCommand extends Logging { try { if (opts.options.has(opts.listOpt)) - consumerGroupService.listGroups().foreach(println(_)) + consumerGroupService.listGroups() else if (opts.options.has(opts.describeOpt)) consumerGroupService.describeGroups() else if (opts.options.has(opts.deleteOpt)) @@ -84,6 +87,10 @@ object ConsumerGroupCommand extends Logging { } } + def consumerGroupStatesFromString(input: String): List[ConsumerGroupState] = { + return input.split(',').map(s => ConsumerGroupState.parse(s.trim.toLowerCase.capitalize)).toSet.toList + } + val MISSING_COLUMN_VALUE = "-" def printError(msg: String, e: Option[Throwable] = None): Unit = { @@ -178,12 +185,44 @@ object ConsumerGroupCommand extends Logging { } else None } - def listGroups(): List[String] = { + def listGroups(): Unit = { + if (opts.options.has(opts.stateOpt)) { + val stateValue = opts.options.valueOf(opts.stateOpt) + val states = if (stateValue == null || stateValue.isEmpty) + allStates + else + consumerGroupStatesFromString(stateValue) + val listings = listConsumerGroupsWithState(states) + printGroupStates(listings.map(e => (e.groupId, e.state.toString)).toList) + } else + listConsumerGroups().foreach(println(_)) + } + + def listConsumerGroups(): List[String] = { val result = adminClient.listConsumerGroups(withTimeoutMs(new ListConsumerGroupsOptions)) val listings = result.all.get.asScala listings.map(_.groupId).toList } + def listConsumerGroupsWithState(states: List[ConsumerGroupState]): List[ConsumerGroupListing] = { + val listConsumerGroupsOptions = withTimeoutMs(new ListConsumerGroupsOptions()) + listConsumerGroupsOptions.inStates(new java.util.HashSet(states.asJava)) + val result = adminClient.listConsumerGroups(listConsumerGroupsOptions) + result.all.get.asScala.toList + } + + private def printGroupStates(groupsAndStates: List[(String, String)]): Unit = { + // find proper columns width + var maxGroupLen = 15 + for ((groupId, state) <- groupsAndStates) { + maxGroupLen = Math.max(maxGroupLen, groupId.length) + } + println(s"%${-maxGroupLen}s %s".format("GROUP", "STATE")) + for ((groupId, state) <- groupsAndStates) { + println(s"%${-maxGroupLen}s %s".format(groupId, state)) + } + } + private def shouldPrintMemberState(group: String, state: Option[String], numRows: Option[Int]): Boolean = { // numRows contains the number of data rows, if any, compiled from the API call in the caller method. // if it's undefined or 0, there is no relevant group information to display. @@ -304,7 +343,7 @@ object ConsumerGroupCommand extends Logging { def describeGroups(): Unit = { val groupIds = - if (opts.options.has(opts.allGroupsOpt)) listGroups() + if (opts.options.has(opts.allGroupsOpt)) listConsumerGroups() else opts.options.valuesOf(opts.groupOpt).asScala val membersOptPresent = opts.options.has(opts.membersOpt) val stateOptPresent = opts.options.has(opts.stateOpt) @@ -368,7 +407,7 @@ object ConsumerGroupCommand extends Logging { def resetOffsets(): Map[String, Map[TopicPartition, OffsetAndMetadata]] = { val groupIds = - if (opts.options.has(opts.allGroupsOpt)) listGroups() + if (opts.options.has(opts.allGroupsOpt)) listConsumerGroups() else opts.options.valuesOf(opts.groupOpt).asScala val consumerGroups = adminClient.describeConsumerGroups( @@ -855,7 +894,7 @@ object ConsumerGroupCommand extends Logging { def deleteGroups(): Map[String, Throwable] = { val groupIds = - if (opts.options.has(opts.allGroupsOpt)) listGroups() + if (opts.options.has(opts.allGroupsOpt)) listConsumerGroups() else opts.options.valuesOf(opts.groupOpt).asScala val groupsToDelete = adminClient.deleteConsumerGroups( @@ -939,8 +978,11 @@ object ConsumerGroupCommand extends Logging { val OffsetsDoc = "Describe the group and list all topic partitions in the group along with their offset lag. " + "This is the default sub-action of and may be used with '--describe' and '--bootstrap-server' options only." + nl + "Example: --bootstrap-server localhost:9092 --describe --group group1 --offsets" - val StateDoc = "Describe the group state. This option may be used with '--describe' and '--bootstrap-server' options only." + nl + - "Example: --bootstrap-server localhost:9092 --describe --group group1 --state" + val StateDoc = "When specified with '--describe', includes the state of the group." + nl + + "Example: --bootstrap-server localhost:9092 --describe --group group1 --state" + nl + + "When specified with '--list', it displays the state of all groups. It can also be used to list groups with specific states." + nl + + "Example: --bootstrap-server localhost:9092 --list --state stable,empty" + nl + + "This option may be used with '--describe', '--list' and '--bootstrap-server' options only." val DeleteOffsetsDoc = "Delete offsets of consumer group. Supports one consumer group at the time, and multiple topics." val bootstrapServerOpt = parser.accepts("bootstrap-server", BootstrapServerDoc) @@ -1004,9 +1046,9 @@ object ConsumerGroupCommand extends Logging { val offsetsOpt = parser.accepts("offsets", OffsetsDoc) .availableIf(describeOpt) val stateOpt = parser.accepts("state", StateDoc) - .availableIf(describeOpt) - - parser.mutuallyExclusive(membersOpt, offsetsOpt, stateOpt) + .availableIf(describeOpt, listOpt) + .withOptionalArg() + .ofType(classOf[String]) options = parser.parse(args : _*) @@ -1024,6 +1066,11 @@ object ConsumerGroupCommand extends Logging { if (!options.has(groupOpt) && !options.has(allGroupsOpt)) CommandLineUtils.printUsageAndDie(parser, s"Option $describeOpt takes one of these options: ${allGroupSelectionScopeOpts.mkString(", ")}") + val mutuallyExclusiveOpts: Set[OptionSpec[_]] = Set(membersOpt, offsetsOpt, stateOpt) + if (mutuallyExclusiveOpts.toList.map(o => if (options.has(o)) 1 else 0).sum > 1) { + CommandLineUtils.printUsageAndDie(parser, + s"Option $describeOpt takes at most one of these options: $mutuallyExclusiveOpts") + } } else { if (options.has(timeoutMsOpt)) debug(s"Option $timeoutMsOpt is applicable only when $describeOpt is used.") diff --git a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala index def9e8155550b..13ec302591886 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala @@ -797,12 +797,13 @@ class GroupCoordinator(val brokerId: Int, } } - def handleListGroups(): (Errors, List[GroupOverview]) = { + def handleListGroups(states: List[String]): (Errors, List[GroupOverview]) = { if (!isActive.get) { (Errors.COORDINATOR_NOT_AVAILABLE, List[GroupOverview]()) } else { val errorCode = if (groupManager.isLoading) Errors.COORDINATOR_LOAD_IN_PROGRESS else Errors.NONE - (errorCode, groupManager.currentGroups.map(_.overview).toList) + // if states is empty, return all groups + (errorCode, groupManager.currentGroups.filter(g => states.isEmpty || states.contains(g.summary.state)).map(_.overview).toList) } } diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadata.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadata.scala index 7e3b470b1589b..3aa12929022eb 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadata.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadata.scala @@ -159,7 +159,8 @@ private object GroupMetadata extends Logging { * Case class used to represent group metadata for the ListGroups API */ case class GroupOverview(groupId: String, - protocolType: String) + protocolType: String, + state: GroupState) /** * Case class used to represent group metadata for the DescribeGroup API @@ -562,7 +563,7 @@ private[group] class GroupMetadata(val groupId: String, initialState: GroupState } def overview: GroupOverview = { - GroupOverview(groupId, protocolType.getOrElse("")) + GroupOverview(groupId, protocolType.getOrElse(""), state) } def initializeOffsets(offsets: collection.Map[TopicPartition, CommitRecordMetadataAndOffset], diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 049c4e09c609a..c421ab5cd50a0 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -1397,16 +1397,22 @@ class KafkaApis(val requestChannel: RequestChannel, } def handleListGroupsRequest(request: RequestChannel.Request): Unit = { - val (error, groups) = groupCoordinator.handleListGroups() + val listGroupsRequest = request.body[ListGroupsRequest] + val states = listGroupsRequest.data.states.asScala.toList + val (error, groups) = groupCoordinator.handleListGroups(states) if (authorize(request.context, DESCRIBE, CLUSTER, CLUSTER_NAME)) // With describe cluster access all groups are returned. We keep this alternative for backward compatibility. sendResponseMaybeThrottle(request, requestThrottleMs => new ListGroupsResponse(new ListGroupsResponseData() .setErrorCode(error.code) - .setGroups(groups.map { group => new ListGroupsResponseData.ListedGroup() - .setGroupId(group.groupId) - .setProtocolType(group.protocolType)}.asJava - ) + .setGroups(groups.map { group => + val listedGroup = new ListGroupsResponseData.ListedGroup() + .setGroupId(group.groupId) + .setProtocolType(group.protocolType) + if (!states.isEmpty) + listedGroup.setGroupState(group.state.toString) + listedGroup + }.asJava) .setThrottleTimeMs(requestThrottleMs) )) else { @@ -1414,10 +1420,14 @@ class KafkaApis(val requestChannel: RequestChannel, sendResponseMaybeThrottle(request, requestThrottleMs => new ListGroupsResponse(new ListGroupsResponseData() .setErrorCode(error.code) - .setGroups(filteredGroups.map { group => new ListGroupsResponseData.ListedGroup() - .setGroupId(group.groupId) - .setProtocolType(group.protocolType)}.asJava - ) + .setGroups(filteredGroups.map { group => + val listedGroup = new ListGroupsResponseData.ListedGroup() + .setGroupId(group.groupId) + .setProtocolType(group.protocolType) + if (!states.isEmpty) + listedGroup.setGroupState(group.state.toString) + listedGroup + }.asJava) .setThrottleTimeMs(requestThrottleMs) )) } diff --git a/core/src/test/scala/integration/kafka/api/SaslClientsWithInvalidCredentialsTest.scala b/core/src/test/scala/integration/kafka/api/SaslClientsWithInvalidCredentialsTest.scala index 03c5f48493dd5..8442471e5b8c5 100644 --- a/core/src/test/scala/integration/kafka/api/SaslClientsWithInvalidCredentialsTest.scala +++ b/core/src/test/scala/integration/kafka/api/SaslClientsWithInvalidCredentialsTest.scala @@ -176,7 +176,7 @@ class SaslClientsWithInvalidCredentialsTest extends IntegrationTestHarness with consumer.subscribe(List(topic).asJava) verifyWithRetry(consumer.poll(Duration.ofMillis(1000))) - assertEquals(1, consumerGroupService.listGroups.size) + assertEquals(1, consumerGroupService.listConsumerGroups.size) consumerGroupService.close() } diff --git a/core/src/test/scala/unit/kafka/admin/DeleteConsumerGroupsTest.scala b/core/src/test/scala/unit/kafka/admin/DeleteConsumerGroupsTest.scala index 23de7904d5eb6..a3fa8f511e137 100644 --- a/core/src/test/scala/unit/kafka/admin/DeleteConsumerGroupsTest.scala +++ b/core/src/test/scala/unit/kafka/admin/DeleteConsumerGroupsTest.scala @@ -107,7 +107,7 @@ class DeleteConsumerGroupsTest extends ConsumerGroupCommandTest { val service = getConsumerGroupService(cgcArgs) TestUtils.waitUntilTrue(() => { - service.listGroups().contains(group) && service.collectGroupState(group).state == "Stable" + service.listConsumerGroups().contains(group) && service.collectGroupState(group).state == "Stable" }, "The group did not initialize as expected.") executor.shutdown() @@ -137,7 +137,7 @@ class DeleteConsumerGroupsTest extends ConsumerGroupCommandTest { val service = getConsumerGroupService(cgcArgs) TestUtils.waitUntilTrue(() => { - service.listGroups().toSet == groups.keySet && + service.listConsumerGroups().toSet == groups.keySet && groups.keySet.forall(groupId => service.collectGroupState(groupId).state == "Stable") }, "The group did not initialize as expected.") @@ -169,7 +169,7 @@ class DeleteConsumerGroupsTest extends ConsumerGroupCommandTest { val service = getConsumerGroupService(cgcArgs) TestUtils.waitUntilTrue(() => { - service.listGroups().contains(group) && service.collectGroupState(group).state == "Stable" + service.listConsumerGroups().contains(group) && service.collectGroupState(group).state == "Stable" }, "The group did not initialize as expected.") executor.shutdown() @@ -194,7 +194,7 @@ class DeleteConsumerGroupsTest extends ConsumerGroupCommandTest { val service = getConsumerGroupService(cgcArgs) TestUtils.waitUntilTrue(() => { - service.listGroups().contains(group) && service.collectGroupState(group).state == "Stable" + service.listConsumerGroups().contains(group) && service.collectGroupState(group).state == "Stable" }, "The group did not initialize as expected.") executor.shutdown() @@ -221,7 +221,7 @@ class DeleteConsumerGroupsTest extends ConsumerGroupCommandTest { val service = getConsumerGroupService(cgcArgs) TestUtils.waitUntilTrue(() => { - service.listGroups().contains(group) && service.collectGroupState(group).state == "Stable" + service.listConsumerGroups().contains(group) && service.collectGroupState(group).state == "Stable" }, "The group did not initialize as expected.") executor.shutdown() diff --git a/core/src/test/scala/unit/kafka/admin/DescribeConsumerGroupTest.scala b/core/src/test/scala/unit/kafka/admin/DescribeConsumerGroupTest.scala index 560153f3a34aa..a019df4a605d9 100644 --- a/core/src/test/scala/unit/kafka/admin/DescribeConsumerGroupTest.scala +++ b/core/src/test/scala/unit/kafka/admin/DescribeConsumerGroupTest.scala @@ -18,7 +18,6 @@ package kafka.admin import java.util.Properties -import joptsimple.OptionException import kafka.utils.TestUtils import org.apache.kafka.clients.consumer.{ConsumerConfig, RoundRobinAssignor} import org.apache.kafka.common.TopicPartition @@ -52,13 +51,6 @@ class DescribeConsumerGroupTest extends ConsumerGroupCommandTest { } } - @Test(expected = classOf[OptionException]) - def testDescribeWithMultipleSubActions(): Unit = { - TestUtils.createOffsetsTopic(zkClient, servers) - val cgcArgs = Array("--bootstrap-server", brokerList, "--describe", "--group", group, "--members", "--state") - getConsumerGroupService(cgcArgs) - } - @Test def testDescribeOffsetsOfNonExistingGroup(): Unit = { val group = "missing.group" diff --git a/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala b/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala index 7429a43a1f102..953fa750c80a5 100644 --- a/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala +++ b/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala @@ -17,8 +17,11 @@ package kafka.admin import joptsimple.OptionException +import org.junit.Assert._ import org.junit.Test import kafka.utils.TestUtils +import org.apache.kafka.common.ConsumerGroupState +import org.apache.kafka.clients.admin.ConsumerGroupListing class ListConsumerGroupTest extends ConsumerGroupCommandTest { @@ -34,7 +37,7 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { val expectedGroups = Set(group, simpleGroup) var foundGroups = Set.empty[String] TestUtils.waitUntilTrue(() => { - foundGroups = service.listGroups().toSet + foundGroups = service.listConsumerGroups().toSet expectedGroups == foundGroups }, s"Expected --list to show groups $expectedGroups, but found $foundGroups.") } @@ -44,4 +47,52 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { val cgcArgs = Array("--new-consumer", "--bootstrap-server", brokerList, "--list") getConsumerGroupService(cgcArgs) } + + @Test + def testListConsumerGroupsWithStates(): Unit = { + val simpleGroup = "simple-group" + addSimpleGroupExecutor(group = simpleGroup) + addConsumerGroupExecutor(numConsumers = 1) + + val cgcArgs = Array("--bootstrap-server", brokerList, "--list", "--state") + val service = getConsumerGroupService(cgcArgs) + + val expectedListing = Set( + new ConsumerGroupListing(simpleGroup, true, ConsumerGroupState.EMPTY), + new ConsumerGroupListing(group, false, ConsumerGroupState.STABLE)) + + var foundListing = Set.empty[ConsumerGroupListing] + TestUtils.waitUntilTrue(() => { + foundListing = service.listConsumerGroupsWithState(ConsumerGroupState.values.toList).toSet + expectedListing == foundListing + }, s"Expected to show groups $expectedListing, but found $foundListing") + + val expectedListingStable = Set( + new ConsumerGroupListing(group, false, ConsumerGroupState.STABLE)) + + foundListing = Set.empty[ConsumerGroupListing] + TestUtils.waitUntilTrue(() => { + foundListing = service.listConsumerGroupsWithState(List(ConsumerGroupState.STABLE)).toSet + expectedListingStable == foundListing + }, s"Expected to show groups expectedListingStable, but found $foundListing") + } + + @Test + def testConsumerGroupStatesFromString(): Unit = { + var result = ConsumerGroupCommand.consumerGroupStatesFromString("bad, wrong") + assertEquals(List(ConsumerGroupState.UNKNOWN), result) + + result = ConsumerGroupCommand.consumerGroupStatesFromString(" bad, ") + assertEquals(List(ConsumerGroupState.UNKNOWN), result) + + result = ConsumerGroupCommand.consumerGroupStatesFromString(" bad, stable") + assertEquals(List(ConsumerGroupState.UNKNOWN, ConsumerGroupState.STABLE), result) + + result = ConsumerGroupCommand.consumerGroupStatesFromString("STABLE, stable, Stable, eMpTy") + assertEquals(List(ConsumerGroupState.STABLE, ConsumerGroupState.EMPTY), result) + + result = ConsumerGroupCommand.consumerGroupStatesFromString(" , ,") + assertEquals(List(ConsumerGroupState.UNKNOWN), result) + } + } diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala index 3200f81c9efec..96f154094e6c9 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala @@ -174,7 +174,7 @@ class GroupCoordinatorTest { assertEquals(Errors.COORDINATOR_LOAD_IN_PROGRESS, describeGroupError) // ListGroups - val (listGroupsError, _) = groupCoordinator.handleListGroups() + val (listGroupsError, _) = groupCoordinator.handleListGroups(List()) assertEquals(Errors.COORDINATOR_LOAD_IN_PROGRESS, listGroupsError) // DeleteGroups @@ -3208,10 +3208,10 @@ class GroupCoordinatorTest { val syncGroupResult = syncGroupLeader(groupId, generationId, assignedMemberId, Map(assignedMemberId -> Array[Byte]())) assertEquals(Errors.NONE, syncGroupResult.error) - val (error, groups) = groupCoordinator.handleListGroups() + val (error, groups) = groupCoordinator.handleListGroups(List()) assertEquals(Errors.NONE, error) assertEquals(1, groups.size) - assertEquals(GroupOverview("groupId", "consumer"), groups.head) + assertEquals(GroupOverview("groupId", "consumer", Stable), groups.head) } @Test @@ -3220,10 +3220,56 @@ class GroupCoordinatorTest { val joinGroupResult = dynamicJoinGroup(groupId, memberId, protocolType, protocols) assertEquals(Errors.NONE, joinGroupResult.error) - val (error, groups) = groupCoordinator.handleListGroups() + val (error, groups) = groupCoordinator.handleListGroups(List()) assertEquals(Errors.NONE, error) assertEquals(1, groups.size) - assertEquals(GroupOverview("groupId", "consumer"), groups.head) + assertEquals(GroupOverview("groupId", "consumer", CompletingRebalance), groups.head) + } + + @Test + def testListGroupsWithStates(): Unit = { + val allStates = List(PreparingRebalance, CompletingRebalance, Stable, Dead, Empty).map(s => s.toString) + val memberId = JoinGroupRequest.UNKNOWN_MEMBER_ID + + // Member joins the group + val joinGroupResult = dynamicJoinGroup(groupId, memberId, protocolType, protocols) + val assignedMemberId = joinGroupResult.memberId + val generationId = joinGroupResult.generationId + assertEquals(Errors.NONE, joinGroupResult.error) + + // The group should be in CompletingRebalance + val (error, groups) = groupCoordinator.handleListGroups(List(CompletingRebalance.toString)) + assertEquals(Errors.NONE, error) + assertEquals(1, groups.size) + val (error2, groups2) = groupCoordinator.handleListGroups(allStates.filterNot(s => s == CompletingRebalance.toString)) + assertEquals(Errors.NONE, error2) + assertEquals(0, groups2.size) + + // Member syncs + EasyMock.reset(replicaManager) + val syncGroupResult = syncGroupLeader(groupId, generationId, assignedMemberId, Map(assignedMemberId -> Array[Byte]())) + assertEquals(Errors.NONE, syncGroupResult.error) + + // The group is now stable + val (error3, groups3) = groupCoordinator.handleListGroups(List(Stable.toString)) + assertEquals(Errors.NONE, error3) + assertEquals(1, groups3.size) + val (error4, groups4) = groupCoordinator.handleListGroups(allStates.filterNot(s => s == Stable.toString)) + assertEquals(Errors.NONE, error4) + assertEquals(0, groups4.size) + + // Member leaves + EasyMock.reset(replicaManager) + val leaveGroupResults = singleLeaveGroup(groupId, assignedMemberId) + verifyLeaveGroupResult(leaveGroupResults) + + // The group is now empty + val (error5, groups5) = groupCoordinator.handleListGroups(List(Empty.toString)) + assertEquals(Errors.NONE, error5) + assertEquals(1, groups5.size) + val (error6, groups6) = groupCoordinator.handleListGroups(allStates.filterNot(s => s == Empty.toString)) + assertEquals(Errors.NONE, error6) + assertEquals(0, groups6.size) } @Test From 143514cc2a880e2896085cdeb34cc697d62296ab Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Mon, 16 Mar 2020 17:20:21 +0000 Subject: [PATCH 02/14] Add test in admin e2e --- .../api/PlaintextAdminIntegrationTest.scala | 21 +++++++++++++++++-- 1 file changed, 19 insertions(+), 2 deletions(-) diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index f014481951290..da25086514df3 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1061,10 +1061,27 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { assertTrue(latch.await(30000, TimeUnit.MILLISECONDS)) // Test that we can list the new group. TestUtils.waitUntilTrue(() => { - val matching = client.listConsumerGroups.all.get().asScala.filter(_.groupId == testGroupId) - matching.nonEmpty + val matching = client.listConsumerGroups.all.get.asScala.filter(group => + group.groupId == testGroupId && + group.state == ConsumerGroupState.UNKNOWN) + matching.nonEmpty && matching.size == 1 }, s"Expected to be able to list $testGroupId") + TestUtils.waitUntilTrue(() => { + val options = new ListConsumerGroupsOptions().inAnyState + val matching = client.listConsumerGroups(options).all.get.asScala.filter(group => + group.groupId == testGroupId && + group.state == ConsumerGroupState.STABLE) + matching.nonEmpty && matching.size == 1 + }, s"Expected to be able to list $testGroupId in state Stable") + + TestUtils.waitUntilTrue(() => { + val options = new ListConsumerGroupsOptions().inStates(Set(ConsumerGroupState.EMPTY).asJava) + val matching = client.listConsumerGroups(options).all.get.asScala.filter( + _.groupId == testGroupId) + matching.isEmpty + }, s"Expected to find zero groups") + val describeWithFakeGroupResult = client.describeConsumerGroups(Seq(testGroupId, fakeGroupId).asJava, new DescribeConsumerGroupsOptions().includeAuthorizedOperations(true)) assertEquals(2, describeWithFakeGroupResult.describedGroups().size()) From badbc1a40e59fdb56ba53c3f9c8633a705da085d Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Mon, 23 Mar 2020 14:50:37 +0000 Subject: [PATCH 03/14] First pass at addressing feedback --- .../common/message/ListGroupsRequest.json | 4 +- .../clients/admin/KafkaAdminClientTest.java | 24 +++++----- .../coordinator/group/GroupCoordinator.scala | 10 ++++- .../coordinator/group/GroupMetadata.scala | 4 +- .../main/scala/kafka/server/KafkaApis.scala | 32 ++++++-------- .../api/PlaintextAdminIntegrationTest.scala | 4 +- .../group/GroupCoordinatorTest.scala | 4 +- .../unit/kafka/server/KafkaApisTest.scala | 44 +++++++++++++++++++ 8 files changed, 85 insertions(+), 41 deletions(-) diff --git a/clients/src/main/resources/common/message/ListGroupsRequest.json b/clients/src/main/resources/common/message/ListGroupsRequest.json index 513d822c7e803..ca72a9c04ba5a 100644 --- a/clients/src/main/resources/common/message/ListGroupsRequest.json +++ b/clients/src/main/resources/common/message/ListGroupsRequest.json @@ -26,8 +26,8 @@ "flexibleVersions": "3+", "fields": [ { "name": "States", "type": "[]string", "versions": "4+", "tag": 0, "taggedVersions": "4+", - "about": "The states of the groups we want to list", "fields": [ - { "name": "Name", "type": "string", "versions": "3+", "about": "The group state" } + "about": "The states of the groups we want to list. If empty, all groups are returned.", "fields": [ + { "name": "Name", "type": "string", "versions": "4+", "about": "The group state" } ]} ] } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index 2dc8274ab5d36..720aa21ef992a 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -1322,19 +1322,17 @@ public void testListConsumerGroupsWithStates() throws Exception { env.kafkaClient().prepareResponse(prepareMetadataResponse(env.cluster(), Errors.NONE)); env.kafkaClient().prepareResponseFrom( - new ListGroupsResponse( - new ListGroupsResponseData() - .setErrorCode(Errors.NONE.code()) - .setGroups(Arrays.asList( - new ListGroupsResponseData.ListedGroup() - .setGroupId("group-1") - .setProtocolType(ConsumerProtocol.PROTOCOL_TYPE) - .setGroupState("Stable"), - new ListGroupsResponseData.ListedGroup() - .setGroupId("group-2") - .setGroupState("Empty") - ))), - env.cluster().nodeById(0)); + new ListGroupsResponse(new ListGroupsResponseData() + .setErrorCode(Errors.NONE.code()) + .setGroups(Arrays.asList( + new ListGroupsResponseData.ListedGroup() + .setGroupId("group-1") + .setProtocolType(ConsumerProtocol.PROTOCOL_TYPE) + .setGroupState("Stable"), + new ListGroupsResponseData.ListedGroup() + .setGroupId("group-2") + .setGroupState("Empty")))), + env.cluster().nodeById(0)); final ListConsumerGroupsOptions options = new ListConsumerGroupsOptions().inAnyState(); final ListConsumerGroupsResult result = env.adminClient().listConsumerGroups(options); diff --git a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala index 13ec302591886..a93856cc1f622 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala @@ -802,8 +802,16 @@ class GroupCoordinator(val brokerId: Int, (Errors.COORDINATOR_NOT_AVAILABLE, List[GroupOverview]()) } else { val errorCode = if (groupManager.isLoading) Errors.COORDINATOR_LOAD_IN_PROGRESS else Errors.NONE + //(errorCode, groupManager.currentGroups.filter(g => states.isEmpty || states.contains(g.summary.state)).map(_.overview).toList) + println("3") // if states is empty, return all groups - (errorCode, groupManager.currentGroups.filter(g => states.isEmpty || states.contains(g.summary.state)).map(_.overview).toList) + val groups = if (states.isEmpty) + groupManager.currentGroups + else + groupManager.currentGroups.filter(g => states.contains(g.summary.state)) + println("4") + println(groups) + (errorCode, groups.map(_.overview).toList) } } diff --git a/core/src/main/scala/kafka/coordinator/group/GroupMetadata.scala b/core/src/main/scala/kafka/coordinator/group/GroupMetadata.scala index 3aa12929022eb..f82b6639ada46 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupMetadata.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupMetadata.scala @@ -160,7 +160,7 @@ private object GroupMetadata extends Logging { */ case class GroupOverview(groupId: String, protocolType: String, - state: GroupState) + state: String) /** * Case class used to represent group metadata for the DescribeGroup API @@ -563,7 +563,7 @@ private[group] class GroupMetadata(val groupId: String, initialState: GroupState } def overview: GroupOverview = { - GroupOverview(groupId, protocolType.getOrElse(""), state) + GroupOverview(groupId, protocolType.getOrElse(""), state.toString) } def initializeOffsets(offsets: collection.Map[TopicPartition, CommitRecordMetadataAndOffset], diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index c421ab5cd50a0..3f9660a477757 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -86,6 +86,7 @@ import scala.jdk.CollectionConverters._ import scala.collection.mutable.ArrayBuffer import scala.collection.{Map, Seq, Set, immutable, mutable} import scala.util.{Failure, Success, Try} +import kafka.coordinator.group.GroupOverview /** @@ -1399,11 +1400,9 @@ class KafkaApis(val requestChannel: RequestChannel, def handleListGroupsRequest(request: RequestChannel.Request): Unit = { val listGroupsRequest = request.body[ListGroupsRequest] val states = listGroupsRequest.data.states.asScala.toList - val (error, groups) = groupCoordinator.handleListGroups(states) - if (authorize(request.context, DESCRIBE, CLUSTER, CLUSTER_NAME)) - // With describe cluster access all groups are returned. We keep this alternative for backward compatibility. - sendResponseMaybeThrottle(request, requestThrottleMs => - new ListGroupsResponse(new ListGroupsResponseData() + + def createResponse(throttleMs: Int, groups: List[GroupOverview], error: Errors): AbstractResponse = { + new ListGroupsResponse(new ListGroupsResponseData() .setErrorCode(error.code) .setGroups(groups.map { group => val listedGroup = new ListGroupsResponseData.ListedGroup() @@ -1413,23 +1412,18 @@ class KafkaApis(val requestChannel: RequestChannel, listedGroup.setGroupState(group.state.toString) listedGroup }.asJava) - .setThrottleTimeMs(requestThrottleMs) - )) + .setThrottleTimeMs(throttleMs) + ) + } + val (error, groups) = groupCoordinator.handleListGroups(states) + if (authorize(request.context, DESCRIBE, CLUSTER, CLUSTER_NAME)) + // With describe cluster access all groups are returned. We keep this alternative for backward compatibility. + sendResponseMaybeThrottle(request, requestThrottleMs => + createResponse(requestThrottleMs, groups, error)) else { val filteredGroups = groups.filter(group => authorize(request.context, DESCRIBE, GROUP, group.groupId)) sendResponseMaybeThrottle(request, requestThrottleMs => - new ListGroupsResponse(new ListGroupsResponseData() - .setErrorCode(error.code) - .setGroups(filteredGroups.map { group => - val listedGroup = new ListGroupsResponseData.ListedGroup() - .setGroupId(group.groupId) - .setProtocolType(group.protocolType) - if (!states.isEmpty) - listedGroup.setGroupState(group.state.toString) - listedGroup - }.asJava) - .setThrottleTimeMs(requestThrottleMs) - )) + createResponse(requestThrottleMs, filteredGroups, error)) } } diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index da25086514df3..e63b59ce536c4 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1064,7 +1064,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { val matching = client.listConsumerGroups.all.get.asScala.filter(group => group.groupId == testGroupId && group.state == ConsumerGroupState.UNKNOWN) - matching.nonEmpty && matching.size == 1 + matching.size == 1 }, s"Expected to be able to list $testGroupId") TestUtils.waitUntilTrue(() => { @@ -1072,7 +1072,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { val matching = client.listConsumerGroups(options).all.get.asScala.filter(group => group.groupId == testGroupId && group.state == ConsumerGroupState.STABLE) - matching.nonEmpty && matching.size == 1 + matching.size == 1 }, s"Expected to be able to list $testGroupId in state Stable") TestUtils.waitUntilTrue(() => { diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala index 96f154094e6c9..90eb934ef8bf5 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala @@ -3211,7 +3211,7 @@ class GroupCoordinatorTest { val (error, groups) = groupCoordinator.handleListGroups(List()) assertEquals(Errors.NONE, error) assertEquals(1, groups.size) - assertEquals(GroupOverview("groupId", "consumer", Stable), groups.head) + assertEquals(GroupOverview("groupId", "consumer", Stable.toString), groups.head) } @Test @@ -3223,7 +3223,7 @@ class GroupCoordinatorTest { val (error, groups) = groupCoordinator.handleListGroups(List()) assertEquals(Errors.NONE, error) assertEquals(1, groups.size) - assertEquals(GroupOverview("groupId", "consumer", CompletingRebalance), groups.head) + assertEquals(GroupOverview("groupId", "consumer", CompletingRebalance.toString), groups.head) } @Test diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index 7dcb70f9725ee..b76134cfe2792 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -28,6 +28,7 @@ import kafka.api.LeaderAndIsr import kafka.api.{ApiVersion, KAFKA_0_10_2_IV0, KAFKA_2_2_IV1} import kafka.cluster.Partition import kafka.controller.KafkaController +import kafka.coordinator.group.{GroupOverview, GroupState} import kafka.coordinator.group.GroupCoordinatorConcurrencyTest.JoinGroupCallback import kafka.coordinator.group.GroupCoordinatorConcurrencyTest.SyncGroupCallback import kafka.coordinator.group.JoinGroupResult @@ -1724,6 +1725,49 @@ class KafkaApisTest { EasyMock.verify(replicaManager) } + def testListGroupsRequest(): Unit = { + val overviews = List( + GroupOverview("group1", "protocol1", "Stable"), + GroupOverview("goupp2", "qwerty", "Empty") + ) + val response = listGroupRequest(Option.empty, overviews) + assertEquals(2, response.data.groups.size) + assertEquals("", response.data.groups.get(0).groupState) + assertEquals("", response.data.groups.get(1).groupState) + } + + @Test + def testListGroupsRequestWithState(): Unit = { + val overviews = List( + GroupOverview("group1", "protocol1", "Stable") + ) + val response = listGroupRequest(Option.apply("Stable"), overviews) + assertEquals(1, response.data.groups.size) + assertEquals("Stable", response.data.groups.get(0).groupState) + } + + private def listGroupRequest(state: Option[String], overviews: List[GroupOverview]): ListGroupsResponse = { + EasyMock.reset(groupCoordinator, clientRequestQuotaManager, requestChannel) + + val data = new ListGroupsRequestData() + if (state.isDefined) + data.setStates(Collections.singletonList(state.get)) + val listGroupsRequest = new ListGroupsRequest.Builder(data).build() + val requestChannelRequest = buildRequest(listGroupsRequest) + + val capturedResponse = expectNoThrottling() + val expectedStates = if (state.isDefined) List(state.get) else List() + EasyMock.expect(groupCoordinator.handleListGroups(expectedStates)) + .andReturn((Errors.NONE, overviews)) + EasyMock.replay(groupCoordinator, clientRequestQuotaManager, requestChannel) + + createKafkaApis().handleListGroupsRequest(requestChannelRequest) + + val response = readResponse(ApiKeys.LIST_GROUPS, listGroupsRequest, capturedResponse).asInstanceOf[ListGroupsResponse] + assertEquals(Errors.NONE.code, response.data.errorCode) + return response + } + /** * Return pair of listener names in the metadataCache: PLAINTEXT and LISTENER2 respectively. */ From b243a61df164637072453ec1a32ea1cdbfdb4ba1 Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Thu, 26 Mar 2020 14:42:29 +0000 Subject: [PATCH 04/14] Addressed comments --- .../clients/admin/ConsumerGroupListing.java | 30 ++++++++----- .../kafka/clients/admin/KafkaAdminClient.java | 5 ++- .../kafka/common/ConsumerGroupState.java | 1 - .../clients/admin/KafkaAdminClientTest.java | 6 +-- .../kafka/admin/ConsumerGroupCommand.scala | 7 ++- .../coordinator/group/GroupCoordinator.scala | 4 -- .../api/PlaintextAdminIntegrationTest.scala | 4 +- .../kafka/admin/ListConsumerGroupTest.scala | 43 ++++++++++++------- .../unit/kafka/server/KafkaApisTest.scala | 2 +- .../kafka/server/ListGroupsRequestTest.scala | 42 ++++++++++++++++++ 10 files changed, 105 insertions(+), 39 deletions(-) create mode 100644 core/src/test/scala/unit/kafka/server/ListGroupsRequestTest.scala diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java b/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java index 205e7ebacdcfd..40d96129625a0 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java @@ -18,6 +18,7 @@ package org.apache.kafka.clients.admin; import java.util.Objects; +import java.util.Optional; import org.apache.kafka.common.ConsumerGroupState; @@ -27,7 +28,7 @@ public class ConsumerGroupListing { private final String groupId; private final boolean isSimpleConsumerGroup; - private final ConsumerGroupState state; + private final Optional state; /** * Create an instance with the specified parameters. @@ -36,7 +37,7 @@ public class ConsumerGroupListing { * @param isSimpleConsumerGroup If consumer group is simple or not. * @param state The state of the consumer group */ - public ConsumerGroupListing(String groupId, boolean isSimpleConsumerGroup, ConsumerGroupState state) { + public ConsumerGroupListing(String groupId, boolean isSimpleConsumerGroup, Optional state) { this.groupId = groupId; this.isSimpleConsumerGroup = isSimpleConsumerGroup; this.state = state; @@ -59,7 +60,7 @@ public boolean isSimpleConsumerGroup() { /** * Consumer Group state */ - public ConsumerGroupState state() { + public Optional state() { return state; } @@ -78,18 +79,25 @@ public int hashCode() { } @Override - public boolean equals(final Object o) { - if (this == o) return true; - if (o == null || getClass() != o.getClass()) return false; - final ConsumerGroupListing that = (ConsumerGroupListing) o; + public boolean equals(Object obj) { + if (this == obj) + return true; + if (obj == null) + return false; + if (getClass() != obj.getClass()) + return false; + ConsumerGroupListing other = (ConsumerGroupListing) obj; if (groupId == null) { - if (that.groupId != null) + if (other.groupId != null) return false; - } else if (!groupId.equals(that.groupId)) + } else if (!groupId.equals(other.groupId)) return false; - if (isSimpleConsumerGroup != that.isSimpleConsumerGroup) + if (isSimpleConsumerGroup != other.isSimpleConsumerGroup) return false; - if (state != that.state) + if (state == null) { + if (other.state != null) + return false; + } else if (!state.equals(other.state)) return false; return true; } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 0e37acad5770c..2736dde71632a 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3063,7 +3063,10 @@ private void maybeAddConsumerGroup(ListGroupsResponseData.ListedGroup group) { String protocolType = group.protocolType(); if (protocolType.equals(ConsumerProtocol.PROTOCOL_TYPE) || protocolType.isEmpty()) { final String groupId = group.groupId(); - final ConsumerGroupListing groupListing = new ConsumerGroupListing(groupId, protocolType.isEmpty(), ConsumerGroupState.parse(group.groupState())); + final Optional state = (group.groupState().isEmpty()) + ? Optional.empty() + : Optional.of(ConsumerGroupState.parse(group.groupState())); + final ConsumerGroupListing groupListing = new ConsumerGroupListing(groupId, protocolType.isEmpty(), state); results.addListing(groupListing); } } diff --git a/clients/src/main/java/org/apache/kafka/common/ConsumerGroupState.java b/clients/src/main/java/org/apache/kafka/common/ConsumerGroupState.java index 7f3d4f0883b10..36d1a4da36118 100644 --- a/clients/src/main/java/org/apache/kafka/common/ConsumerGroupState.java +++ b/clients/src/main/java/org/apache/kafka/common/ConsumerGroupState.java @@ -45,7 +45,6 @@ public enum ConsumerGroupState { this.name = name; } - /** * Parse a string into a consumer group state. */ diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index 720aa21ef992a..51686ffc41b62 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -1283,7 +1283,7 @@ public void testListConsumerGroups() throws Exception { Set groupIds = new HashSet<>(); for (ConsumerGroupListing listing : listings) { groupIds.add(listing.groupId()); - assertEquals(ConsumerGroupState.UNKNOWN, listing.state()); + assertFalse(listing.state().isPresent()); } assertEquals(Utils.mkSet("group-1", "group-2", "group-3"), groupIds); @@ -1340,8 +1340,8 @@ public void testListConsumerGroupsWithStates() throws Exception { assertEquals(2, listings.size()); List expected = new ArrayList<>(); - expected.add(new ConsumerGroupListing("group-2", true, ConsumerGroupState.EMPTY)); - expected.add(new ConsumerGroupListing("group-1", false, ConsumerGroupState.STABLE)); + expected.add(new ConsumerGroupListing("group-2", true, Optional.of(ConsumerGroupState.EMPTY))); + expected.add(new ConsumerGroupListing("group-1", false, Optional.of(ConsumerGroupState.STABLE))); assertEquals(expected, listings); assertEquals(0, result.errors().get().size()); } diff --git a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala index 1e22f30c9f61a..b12701f0d70af 100755 --- a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala +++ b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala @@ -88,7 +88,12 @@ object ConsumerGroupCommand extends Logging { } def consumerGroupStatesFromString(input: String): List[ConsumerGroupState] = { - return input.split(',').map(s => ConsumerGroupState.parse(s.trim.toLowerCase.capitalize)).toSet.toList + val parsedStates = input.split(',').map(s => ConsumerGroupState.parse(s.trim.toLowerCase.capitalize)).toSet.toList + if (parsedStates.contains(ConsumerGroupState.UNKNOWN)) { + val validStates = ConsumerGroupState.values().filter(_ != ConsumerGroupState.UNKNOWN) + throw new IllegalArgumentException(s"Invalid state list '$input'. Valid states are: ${validStates.mkString(",")}") + } + return parsedStates } val MISSING_COLUMN_VALUE = "-" diff --git a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala index a93856cc1f622..094d3df9a41da 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala @@ -802,15 +802,11 @@ class GroupCoordinator(val brokerId: Int, (Errors.COORDINATOR_NOT_AVAILABLE, List[GroupOverview]()) } else { val errorCode = if (groupManager.isLoading) Errors.COORDINATOR_LOAD_IN_PROGRESS else Errors.NONE - //(errorCode, groupManager.currentGroups.filter(g => states.isEmpty || states.contains(g.summary.state)).map(_.overview).toList) - println("3") // if states is empty, return all groups val groups = if (states.isEmpty) groupManager.currentGroups else groupManager.currentGroups.filter(g => states.contains(g.summary.state)) - println("4") - println(groups) (errorCode, groups.map(_.overview).toList) } } diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index e63b59ce536c4..be5773d5b3578 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1063,7 +1063,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { TestUtils.waitUntilTrue(() => { val matching = client.listConsumerGroups.all.get.asScala.filter(group => group.groupId == testGroupId && - group.state == ConsumerGroupState.UNKNOWN) + !group.state.isPresent) matching.size == 1 }, s"Expected to be able to list $testGroupId") @@ -1071,7 +1071,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { val options = new ListConsumerGroupsOptions().inAnyState val matching = client.listConsumerGroups(options).all.get.asScala.filter(group => group.groupId == testGroupId && - group.state == ConsumerGroupState.STABLE) + group.state.get == ConsumerGroupState.STABLE) matching.size == 1 }, s"Expected to be able to list $testGroupId in state Stable") diff --git a/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala b/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala index 953fa750c80a5..f4ab7a9c8452e 100644 --- a/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala +++ b/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala @@ -22,6 +22,7 @@ import org.junit.Test import kafka.utils.TestUtils import org.apache.kafka.common.ConsumerGroupState import org.apache.kafka.clients.admin.ConsumerGroupListing +import java.util.Optional class ListConsumerGroupTest extends ConsumerGroupCommandTest { @@ -58,8 +59,8 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { val service = getConsumerGroupService(cgcArgs) val expectedListing = Set( - new ConsumerGroupListing(simpleGroup, true, ConsumerGroupState.EMPTY), - new ConsumerGroupListing(group, false, ConsumerGroupState.STABLE)) + new ConsumerGroupListing(simpleGroup, true, Optional.of(ConsumerGroupState.EMPTY)), + new ConsumerGroupListing(group, false, Optional.of(ConsumerGroupState.STABLE))) var foundListing = Set.empty[ConsumerGroupListing] TestUtils.waitUntilTrue(() => { @@ -68,7 +69,7 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { }, s"Expected to show groups $expectedListing, but found $foundListing") val expectedListingStable = Set( - new ConsumerGroupListing(group, false, ConsumerGroupState.STABLE)) + new ConsumerGroupListing(group, false, Optional.of(ConsumerGroupState.STABLE))) foundListing = Set.empty[ConsumerGroupListing] TestUtils.waitUntilTrue(() => { @@ -79,20 +80,32 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { @Test def testConsumerGroupStatesFromString(): Unit = { - var result = ConsumerGroupCommand.consumerGroupStatesFromString("bad, wrong") - assertEquals(List(ConsumerGroupState.UNKNOWN), result) - - result = ConsumerGroupCommand.consumerGroupStatesFromString(" bad, ") - assertEquals(List(ConsumerGroupState.UNKNOWN), result) - - result = ConsumerGroupCommand.consumerGroupStatesFromString(" bad, stable") - assertEquals(List(ConsumerGroupState.UNKNOWN, ConsumerGroupState.STABLE), result) - - result = ConsumerGroupCommand.consumerGroupStatesFromString("STABLE, stable, Stable, eMpTy") + val result = ConsumerGroupCommand.consumerGroupStatesFromString("STABLE, stable, Stable, eMpTy") assertEquals(List(ConsumerGroupState.STABLE, ConsumerGroupState.EMPTY), result) - result = ConsumerGroupCommand.consumerGroupStatesFromString(" , ,") - assertEquals(List(ConsumerGroupState.UNKNOWN), result) + try { + ConsumerGroupCommand.consumerGroupStatesFromString("bad, wrong") + } catch { + case e: IllegalArgumentException => //Expected + } + + try { + ConsumerGroupCommand.consumerGroupStatesFromString(" bad, ") + } catch { + case e: IllegalArgumentException => //Expected + } + + try { + ConsumerGroupCommand.consumerGroupStatesFromString(" bad, stable") + } catch { + case e: IllegalArgumentException => //Expected + } + + try { + ConsumerGroupCommand.consumerGroupStatesFromString(" , ,") + } catch { + case e: IllegalArgumentException => //Expected + } } } diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index b76134cfe2792..9202589ecfc30 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -28,7 +28,7 @@ import kafka.api.LeaderAndIsr import kafka.api.{ApiVersion, KAFKA_0_10_2_IV0, KAFKA_2_2_IV1} import kafka.cluster.Partition import kafka.controller.KafkaController -import kafka.coordinator.group.{GroupOverview, GroupState} +import kafka.coordinator.group.GroupOverview import kafka.coordinator.group.GroupCoordinatorConcurrencyTest.JoinGroupCallback import kafka.coordinator.group.GroupCoordinatorConcurrencyTest.SyncGroupCallback import kafka.coordinator.group.JoinGroupResult diff --git a/core/src/test/scala/unit/kafka/server/ListGroupsRequestTest.scala b/core/src/test/scala/unit/kafka/server/ListGroupsRequestTest.scala new file mode 100644 index 0000000000000..d79e984e153d1 --- /dev/null +++ b/core/src/test/scala/unit/kafka/server/ListGroupsRequestTest.scala @@ -0,0 +1,42 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package kafka.server + +import org.apache.kafka.common.errors.UnsupportedVersionException +import org.apache.kafka.common.message.ListGroupsRequestData +import org.apache.kafka.common.requests.ListGroupsRequest +import org.apache.kafka.common.requests.ListGroupsResponse + +import java.util.Arrays + +import org.junit.Assert._ +import org.junit.Test + +class ListGroupsRequestTest extends BaseRequestTest { + override val brokerCount: Int = 1 + + @Test(expected=classOf[UnsupportedVersionException]) + def testListGroupsRequestOldVersionWithStates(): Unit = { + val request = new ListGroupsRequest.Builder( + new ListGroupsRequestData().setStates(Arrays.asList("Stable")) + ).build(3) + + connectAndReceive[ListGroupsResponse](request, destination = controllerSocketServer) + + } + +} From 145b6897390097d93dfedc52036ac2e48906e96a Mon Sep 17 00:00:00 2001 From: Edoardo Comar Date: Thu, 2 Apr 2020 18:31:07 +0100 Subject: [PATCH 05/14] addressed a couple of comments and CLI's list state formatting --- .../java/org/apache/kafka/clients/admin/KafkaAdminClient.java | 2 +- .../src/main/resources/common/message/ListGroupsRequest.json | 2 +- core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 2736dde71632a..0de9036f242ee 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3063,7 +3063,7 @@ private void maybeAddConsumerGroup(ListGroupsResponseData.ListedGroup group) { String protocolType = group.protocolType(); if (protocolType.equals(ConsumerProtocol.PROTOCOL_TYPE) || protocolType.isEmpty()) { final String groupId = group.groupId(); - final Optional state = (group.groupState().isEmpty()) + final Optional state = group.groupState().isEmpty() ? Optional.empty() : Optional.of(ConsumerGroupState.parse(group.groupState())); final ConsumerGroupListing groupListing = new ConsumerGroupListing(groupId, protocolType.isEmpty(), state); diff --git a/clients/src/main/resources/common/message/ListGroupsRequest.json b/clients/src/main/resources/common/message/ListGroupsRequest.json index ca72a9c04ba5a..24f9db02d322e 100644 --- a/clients/src/main/resources/common/message/ListGroupsRequest.json +++ b/clients/src/main/resources/common/message/ListGroupsRequest.json @@ -26,7 +26,7 @@ "flexibleVersions": "3+", "fields": [ { "name": "States", "type": "[]string", "versions": "4+", "tag": 0, "taggedVersions": "4+", - "about": "The states of the groups we want to list. If empty, all groups are returned.", "fields": [ + "about": "The states of the groups we want to list. If empty, all groups are returned, without state.", "fields": [ { "name": "Name", "type": "string", "versions": "4+", "about": "The group state" } ]} ] diff --git a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala index b12701f0d70af..c7db2867910d5 100755 --- a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala +++ b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala @@ -198,7 +198,7 @@ object ConsumerGroupCommand extends Logging { else consumerGroupStatesFromString(stateValue) val listings = listConsumerGroupsWithState(states) - printGroupStates(listings.map(e => (e.groupId, e.state.toString)).toList) + printGroupStates(listings.map(e => (e.groupId, e.state.get.toString)).toList) } else listConsumerGroups().foreach(println(_)) } From a0247b048d6a7f70fb225c82557bec300c69a6a7 Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Thu, 16 Apr 2020 11:33:04 +0100 Subject: [PATCH 06/14] Address latest feedback --- .../admin/ListConsumerGroupsOptions.java | 4 +- .../common/requests/ListGroupsRequest.java | 5 +++ .../common/requests/RequestResponseTest.java | 23 +++++++--- .../kafka/admin/ListConsumerGroupTest.scala | 26 ++++++++++++ .../kafka/server/ListGroupsRequestTest.scala | 42 ------------------- 5 files changed, 50 insertions(+), 50 deletions(-) delete mode 100644 core/src/test/scala/unit/kafka/server/ListGroupsRequestTest.scala diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java index 015a7add496ee..72e6af909723b 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java @@ -46,7 +46,7 @@ public ListConsumerGroupsOptions inStates(Set states) { this.states = Optional.of(states); return this; } - + /** * All groups with their states will be returned by listConsumerGroups() */ @@ -54,7 +54,7 @@ public ListConsumerGroupsOptions inAnyState() { this.states = Optional.of(EnumSet.allOf(ConsumerGroupState.class)); return this; } - + /** * Returns the list of States that are requested */ diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java index 1b48aae6d1f8a..a878e189742ac 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.common.requests; +import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.message.ListGroupsRequestData; import org.apache.kafka.common.message.ListGroupsResponseData; import org.apache.kafka.common.protocol.ApiKeys; @@ -45,6 +46,10 @@ public Builder(ListGroupsRequestData data) { @Override public ListGroupsRequest build(short version) { + if (!data.states().isEmpty() && version < 4) { + throw new UnsupportedVersionException("The broker only supports ListGroups " + + "v" + version + ", but we need v4 or newer to request groups by states."); + } return new ListGroupsRequest(data, version); } diff --git a/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java b/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java index 69b658a16bd99..a591dbb07c8b8 100644 --- a/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java +++ b/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.common.requests; +import org.apache.kafka.common.ConsumerGroupState; import org.apache.kafka.common.ElectionType; import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.Node; @@ -229,8 +230,10 @@ public void testSerialization() throws Exception { checkRequest(createLeaveGroupRequest(), true); checkErrorResponse(createLeaveGroupRequest(), new UnknownServerException(), true); checkResponse(createLeaveGroupResponse(), 0, true); - checkRequest(createListGroupsRequest(), true); - checkErrorResponse(createListGroupsRequest(), new UnknownServerException(), true); + for (short v = ApiKeys.LIST_GROUPS.oldestVersion(); v <= ApiKeys.LIST_GROUPS.latestVersion(); v++) { + checkRequest(createListGroupsRequest(v), false); + checkErrorResponse(createListGroupsRequest(v), new UnknownServerException(), true); + } checkResponse(createListGroupsResponse(), 0, true); checkRequest(createDescribeGroupRequest(), true); checkErrorResponse(createDescribeGroupRequest(), new UnknownServerException(), true); @@ -810,6 +813,13 @@ public void testValidApiVersionsRequest() { assertTrue(request.isValid()); } + @Test(expected = UnsupportedVersionException.class) + public void testListGroupRequestV3FailsWithStates() { + ListGroupsRequestData data = new ListGroupsRequestData() + .setStates(asList(ConsumerGroupState.STABLE.name())); + new ListGroupsRequest.Builder(data).build((short) 3); + } + @Test public void testInvalidApiVersionsRequest() { testInvalidCase("java@apache_kafka", "0.0.0-SNAPSHOT"); @@ -1066,8 +1076,11 @@ private SyncGroupResponse createSyncGroupResponse(int version) { return new SyncGroupResponse(data); } - private ListGroupsRequest createListGroupsRequest() { - return new ListGroupsRequest.Builder(new ListGroupsRequestData()).build(); + private ListGroupsRequest createListGroupsRequest(short version) { + ListGroupsRequestData data = new ListGroupsRequestData(); + if (version >= 4) + data.setStates(Arrays.asList("Stable")); + return new ListGroupsRequest.Builder(data).build(version); } private ListGroupsResponse createListGroupsResponse() { @@ -1133,7 +1146,6 @@ private DeleteGroupsResponse createDeleteGroupsResponse() { ); } - @SuppressWarnings("deprecation") private ListOffsetRequest createListOffsetRequest(int version) { if (version == 0) { Map offsetData = Collections.singletonMap( @@ -1164,7 +1176,6 @@ private ListOffsetRequest createListOffsetRequest(int version) { } } - @SuppressWarnings("deprecation") private ListOffsetResponse createListOffsetResponse(int version) { if (version == 0) { Map responseData = new HashMap<>(); diff --git a/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala b/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala index f4ab7a9c8452e..cd1a31b3191fa 100644 --- a/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala +++ b/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala @@ -108,4 +108,30 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { } } + @Test + def testListGroupCommand(): Unit = { + val simpleGroup = "simple-group" + addSimpleGroupExecutor(group = simpleGroup) + addConsumerGroupExecutor(numConsumers = 1) + var out = "" + + var cgcArgs = Array("--bootstrap-server", brokerList, "--list") + TestUtils.waitUntilTrue(() => { + out = TestUtils.grabConsoleOutput(ConsumerGroupCommand.main(cgcArgs)) + !out.contains("STATE") && out.contains(simpleGroup) && out.contains(group) + }, s"Expected to find $simpleGroup, $group and no header, but found $out") + + cgcArgs = Array("--bootstrap-server", brokerList, "--list", "--state") + TestUtils.waitUntilTrue(() => { + out = TestUtils.grabConsoleOutput(ConsumerGroupCommand.main(cgcArgs)) + out.contains("STATE") && out.contains(simpleGroup) && out.contains(group) + }, s"Expected to find $simpleGroup, $group and the header, but found $out") + + cgcArgs = Array("--bootstrap-server", brokerList, "--list", "--state", "stable") + TestUtils.waitUntilTrue(() => { + out = TestUtils.grabConsoleOutput(ConsumerGroupCommand.main(cgcArgs)) + out.contains("STATE") && out.contains(group) && out.contains("Stable") + }, s"Expected to find $group in state Stable and the header, but found $out") + } + } diff --git a/core/src/test/scala/unit/kafka/server/ListGroupsRequestTest.scala b/core/src/test/scala/unit/kafka/server/ListGroupsRequestTest.scala deleted file mode 100644 index d79e984e153d1..0000000000000 --- a/core/src/test/scala/unit/kafka/server/ListGroupsRequestTest.scala +++ /dev/null @@ -1,42 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package kafka.server - -import org.apache.kafka.common.errors.UnsupportedVersionException -import org.apache.kafka.common.message.ListGroupsRequestData -import org.apache.kafka.common.requests.ListGroupsRequest -import org.apache.kafka.common.requests.ListGroupsResponse - -import java.util.Arrays - -import org.junit.Assert._ -import org.junit.Test - -class ListGroupsRequestTest extends BaseRequestTest { - override val brokerCount: Int = 1 - - @Test(expected=classOf[UnsupportedVersionException]) - def testListGroupsRequestOldVersionWithStates(): Unit = { - val request = new ListGroupsRequest.Builder( - new ListGroupsRequestData().setStates(Arrays.asList("Stable")) - ).build(3) - - connectAndReceive[ListGroupsResponse](request, destination = controllerSocketServer) - - } - -} From 4bed4d3b2a2f45dc9993bb511bfb6deeb3330585 Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Thu, 16 Apr 2020 13:37:42 +0100 Subject: [PATCH 07/14] testDescribeWithMultipleSubActions --- .../admin/DescribeConsumerGroupTest.scala | 23 ++++++++++++++++++- 1 file changed, 22 insertions(+), 1 deletion(-) diff --git a/core/src/test/scala/unit/kafka/admin/DescribeConsumerGroupTest.scala b/core/src/test/scala/unit/kafka/admin/DescribeConsumerGroupTest.scala index a019df4a605d9..79621a7d6abce 100644 --- a/core/src/test/scala/unit/kafka/admin/DescribeConsumerGroupTest.scala +++ b/core/src/test/scala/unit/kafka/admin/DescribeConsumerGroupTest.scala @@ -18,7 +18,7 @@ package kafka.admin import java.util.Properties -import kafka.utils.TestUtils +import kafka.utils.{Exit, TestUtils} import org.apache.kafka.clients.consumer.{ConsumerConfig, RoundRobinAssignor} import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.errors.TimeoutException @@ -51,6 +51,27 @@ class DescribeConsumerGroupTest extends ConsumerGroupCommandTest { } } + @Test + def testDescribeWithMultipleSubActions(): Unit = { + var exitStatus: Option[Int] = None + var exitMessage: Option[String] = None + Exit.setExitProcedure { (status, err) => + exitStatus = Some(status) + exitMessage = err + throw new RuntimeException + } + val cgcArgs = Array("--bootstrap-server", brokerList, "--describe", "--group", group, "--members", "--state") + try { + ConsumerGroupCommand.main(cgcArgs) + } catch { + case e: RuntimeException => //expected + } finally { + Exit.resetExitProcedure() + } + assertEquals(Some(1), exitStatus) + assertTrue(exitMessage.get.contains("Option [describe] takes at most one of these options")) + } + @Test def testDescribeOffsetsOfNonExistingGroup(): Unit = { val group = "missing.group" From d6102e8f3025f70004bb3081785708af76f2855e Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Thu, 23 Apr 2020 13:54:08 +0100 Subject: [PATCH 08/14] Address feedback --- .../apache/kafka/clients/admin/KafkaAdminClient.java | 2 +- .../resources/common/message/ListGroupsResponse.json | 2 +- .../main/scala/kafka/admin/ConsumerGroupCommand.scala | 8 ++++---- .../test/scala/unit/kafka/server/KafkaApisTest.scala | 10 +++++----- 4 files changed, 11 insertions(+), 11 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 0de9036f242ee..d29b2e926156f 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3063,7 +3063,7 @@ private void maybeAddConsumerGroup(ListGroupsResponseData.ListedGroup group) { String protocolType = group.protocolType(); if (protocolType.equals(ConsumerProtocol.PROTOCOL_TYPE) || protocolType.isEmpty()) { final String groupId = group.groupId(); - final Optional state = group.groupState().isEmpty() + final Optional state = group.groupState() == null ? Optional.empty() : Optional.of(ConsumerGroupState.parse(group.groupState())); final ConsumerGroupListing groupListing = new ConsumerGroupListing(groupId, protocolType.isEmpty(), state); diff --git a/clients/src/main/resources/common/message/ListGroupsResponse.json b/clients/src/main/resources/common/message/ListGroupsResponse.json index eb8914edbc42f..dcc96b6fb6f4f 100644 --- a/clients/src/main/resources/common/message/ListGroupsResponse.json +++ b/clients/src/main/resources/common/message/ListGroupsResponse.json @@ -38,7 +38,7 @@ { "name": "ProtocolType", "type": "string", "versions": "0+", "about": "The group protocol type." }, { "name": "GroupState", "type": "string", "versions": "4+", "tag": 0, "taggedVersions": "4+", "ignorable": true, - "about": "The group state string." } + "nullableVersions": "4+", "default": "null", "about": "The group state string." } ]} ] } diff --git a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala index c7db2867910d5..b42d9b130e190 100755 --- a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala +++ b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala @@ -91,9 +91,9 @@ object ConsumerGroupCommand extends Logging { val parsedStates = input.split(',').map(s => ConsumerGroupState.parse(s.trim.toLowerCase.capitalize)).toSet.toList if (parsedStates.contains(ConsumerGroupState.UNKNOWN)) { val validStates = ConsumerGroupState.values().filter(_ != ConsumerGroupState.UNKNOWN) - throw new IllegalArgumentException(s"Invalid state list '$input'. Valid states are: ${validStates.mkString(",")}") + throw new IllegalArgumentException(s"Invalid state list '$input'. Valid states are: ${validStates.mkString(", ")}") } - return parsedStates + parsedStates } val MISSING_COLUMN_VALUE = "-" @@ -198,7 +198,7 @@ object ConsumerGroupCommand extends Logging { else consumerGroupStatesFromString(stateValue) val listings = listConsumerGroupsWithState(states) - printGroupStates(listings.map(e => (e.groupId, e.state.get.toString)).toList) + printGroupStates(listings.map(e => (e.groupId, e.state.get.toString))) } else listConsumerGroups().foreach(println(_)) } @@ -1074,7 +1074,7 @@ object ConsumerGroupCommand extends Logging { val mutuallyExclusiveOpts: Set[OptionSpec[_]] = Set(membersOpt, offsetsOpt, stateOpt) if (mutuallyExclusiveOpts.toList.map(o => if (options.has(o)) 1 else 0).sum > 1) { CommandLineUtils.printUsageAndDie(parser, - s"Option $describeOpt takes at most one of these options: $mutuallyExclusiveOpts") + s"Option $describeOpt takes at most one of these options: ${mutuallyExclusiveOpts.mkString(", ")}") } } else { if (options.has(timeoutMsOpt)) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index 9202589ecfc30..5a5d3edf68bd9 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -1730,10 +1730,10 @@ class KafkaApisTest { GroupOverview("group1", "protocol1", "Stable"), GroupOverview("goupp2", "qwerty", "Empty") ) - val response = listGroupRequest(Option.empty, overviews) + val response = listGroupRequest(None, overviews) assertEquals(2, response.data.groups.size) - assertEquals("", response.data.groups.get(0).groupState) - assertEquals("", response.data.groups.get(1).groupState) + assertNull(response.data.groups.get(0).groupState) + assertNull(response.data.groups.get(1).groupState) } @Test @@ -1741,7 +1741,7 @@ class KafkaApisTest { val overviews = List( GroupOverview("group1", "protocol1", "Stable") ) - val response = listGroupRequest(Option.apply("Stable"), overviews) + val response = listGroupRequest(Some("Stable"), overviews) assertEquals(1, response.data.groups.size) assertEquals("Stable", response.data.groups.get(0).groupState) } @@ -1765,7 +1765,7 @@ class KafkaApisTest { val response = readResponse(ApiKeys.LIST_GROUPS, listGroupsRequest, capturedResponse).asInstanceOf[ListGroupsResponse] assertEquals(Errors.NONE.code, response.data.errorCode) - return response + response } /** From 5854e560da240c68f48b44996f085d37aca9d254 Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Thu, 30 Apr 2020 21:06:55 +0100 Subject: [PATCH 09/14] Improve error reporting --- .../kafka/admin/ConsumerGroupCommand.scala | 26 +++++++++++++------ 1 file changed, 18 insertions(+), 8 deletions(-) diff --git a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala index b42d9b130e190..13cbde94c614f 100755 --- a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala +++ b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala @@ -42,25 +42,33 @@ import scala.collection.immutable.TreeMap import scala.reflect.ClassTag import org.apache.kafka.common.requests.ListOffsetResponse import org.apache.kafka.common.ConsumerGroupState +import joptsimple.OptionException object ConsumerGroupCommand extends Logging { val allStates = ConsumerGroupState.values.toList def main(args: Array[String]): Unit = { - val opts = new ConsumerGroupCommandOptions(args) - CommandLineUtils.printHelpAndExitIfNeeded(opts, "This tool helps to list all consumer groups, describe a consumer group, delete consumer group info, or reset consumer group offsets.") + val opts = new ConsumerGroupCommandOptions(args) + try { + opts.checkArgs() + CommandLineUtils.printHelpAndExitIfNeeded(opts, "This tool helps to list all consumer groups, describe a consumer group, delete consumer group info, or reset consumer group offsets.") - // should have exactly one action - val actions = Seq(opts.listOpt, opts.describeOpt, opts.deleteOpt, opts.resetOffsetsOpt, opts.deleteOffsetsOpt).count(opts.options.has) - if (actions != 1) - CommandLineUtils.printUsageAndDie(opts.parser, "Command must include exactly one action: --list, --describe, --delete, --reset-offsets, --delete-offsets") + // should have exactly one action + val actions = Seq(opts.listOpt, opts.describeOpt, opts.deleteOpt, opts.resetOffsetsOpt, opts.deleteOffsetsOpt).count(opts.options.has) + if (actions != 1) + CommandLineUtils.printUsageAndDie(opts.parser, "Command must include exactly one action: --list, --describe, --delete, --reset-offsets, --delete-offsets") - opts.checkArgs() + run(opts) + } catch { + case e: OptionException => + CommandLineUtils.printUsageAndDie(opts.parser, e.getMessage) + } + } + def run(opts: ConsumerGroupCommandOptions): Unit = { val consumerGroupService = new ConsumerGroupService(opts) - try { if (opts.options.has(opts.listOpt)) consumerGroupService.listGroups() @@ -80,6 +88,8 @@ object ConsumerGroupCommand extends Logging { consumerGroupService.deleteOffsets() } } catch { + case e: IllegalArgumentException => + CommandLineUtils.printUsageAndDie(opts.parser, e.getMessage) case e: Throwable => printError(s"Executing consumer group command failed due to ${e.getMessage}", Some(e)) } finally { From 1f032e71c50411ecff2d99f86e523e02dd560c32 Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Thu, 21 May 2020 15:00:11 +0100 Subject: [PATCH 10/14] Addressed feedback - Switch to use regular fields instead of tagged fields - Empty/null state filter now means all states --- checkstyle/suppressions.xml | 2 +- .../kafka/clients/admin/KafkaAdminClient.java | 12 ++-- .../admin/ListConsumerGroupsOptions.java | 20 +++---- .../common/requests/ListGroupsRequest.java | 2 +- .../common/message/ListGroupsRequest.json | 8 +-- .../common/message/ListGroupsResponse.json | 6 +- .../org/apache/kafka/clients/MockClient.java | 19 ++++-- .../clients/admin/KafkaAdminClientTest.java | 60 +++++++++++++++++-- .../common/requests/RequestResponseTest.java | 27 +++++---- .../kafka/admin/ConsumerGroupCommand.scala | 2 +- .../main/scala/kafka/server/KafkaApis.scala | 9 ++- .../api/PlaintextAdminIntegrationTest.scala | 2 +- .../kafka/admin/ListConsumerGroupTest.scala | 18 ++++-- .../unit/kafka/server/KafkaApisTest.scala | 9 +-- 14 files changed, 131 insertions(+), 65 deletions(-) diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml index f3c463457c01e..3d74333ad9ac9 100644 --- a/checkstyle/suppressions.xml +++ b/checkstyle/suppressions.xml @@ -89,7 +89,7 @@ files="MockAdminClient.java"/> + files="RequestResponseTest.java|FetcherTest.java|KafkaAdminClientTest.java"/> diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index d29b2e926156f..471626cb1dabc 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3052,11 +3052,13 @@ void handleResponse(AbstractResponse abstractResponse) { runnable.call(new Call("listConsumerGroups", deadline, new ConstantNodeIdProvider(node.id())) { @Override ListGroupsRequest.Builder createRequest(int timeoutMs) { - List states = options.states().orElse(Collections.emptySet()) - .stream() - .map(s -> s.toString()) - .collect(Collectors.toList()); - return new ListGroupsRequest.Builder(new ListGroupsRequestData().setStates(states)); + List states = (options.states().isEmpty()) + ? null + : options.states() + .stream() + .map(s -> s.toString()) + .collect(Collectors.toList()); + return new ListGroupsRequest.Builder(new ListGroupsRequestData().setStatesFilter(states)); } private void maybeAddConsumerGroup(ListGroupsResponseData.ListedGroup group) { diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java index 72e6af909723b..81fffe06c6714 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java @@ -17,8 +17,8 @@ package org.apache.kafka.clients.admin; -import java.util.EnumSet; -import java.util.Optional; +import java.util.Collections; +import java.util.HashSet; import java.util.Set; import org.apache.kafka.common.ConsumerGroupState; @@ -32,18 +32,14 @@ @InterfaceStability.Evolving public class ListConsumerGroupsOptions extends AbstractOptions { - private Optional> states = Optional.empty(); + private Set states = Collections.emptySet(); /** * Only groups in these states will be returned by listConsumerGroups() - * If not set, all groups are returned without their states - * throw IllegalArgumentException if states is empty + * If not set, all groups are returned with their states */ public ListConsumerGroupsOptions inStates(Set states) { - if (states == null || states.isEmpty()) { - throw new IllegalArgumentException("states should not be null or empty"); - } - this.states = Optional.of(states); + this.states = (states == null) ? Collections.emptySet() : new HashSet<>(states); return this; } @@ -51,14 +47,14 @@ public ListConsumerGroupsOptions inStates(Set states) { * All groups with their states will be returned by listConsumerGroups() */ public ListConsumerGroupsOptions inAnyState() { - this.states = Optional.of(EnumSet.allOf(ConsumerGroupState.class)); + this.states = Collections.emptySet(); return this; } /** - * Returns the list of States that are requested + * Returns the list of States that are requested or empty if no states have been specified */ - public Optional> states() { + public Set states() { return states; } } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java index a878e189742ac..6ad4e9ce641e7 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java @@ -46,7 +46,7 @@ public Builder(ListGroupsRequestData data) { @Override public ListGroupsRequest build(short version) { - if (!data.states().isEmpty() && version < 4) { + if (data.statesFilter() != null && version < 4) { throw new UnsupportedVersionException("The broker only supports ListGroups " + "v" + version + ", but we need v4 or newer to request groups by states."); } diff --git a/clients/src/main/resources/common/message/ListGroupsRequest.json b/clients/src/main/resources/common/message/ListGroupsRequest.json index 24f9db02d322e..44e7423a0ec71 100644 --- a/clients/src/main/resources/common/message/ListGroupsRequest.json +++ b/clients/src/main/resources/common/message/ListGroupsRequest.json @@ -21,13 +21,13 @@ // // Version 3 is the first flexible version. // - // Version 4 adds the States flexible field (KIP-518). + // Version 4 adds the StatesFilter field (KIP-518). "validVersions": "0-4", "flexibleVersions": "3+", "fields": [ - { "name": "States", "type": "[]string", "versions": "4+", "tag": 0, "taggedVersions": "4+", - "about": "The states of the groups we want to list. If empty, all groups are returned, without state.", "fields": [ - { "name": "Name", "type": "string", "versions": "4+", "about": "The group state" } + { "name": "StatesFilter", "type": "[]string", "versions": "4+", "nullableVersions": "4+", "default": "null", + "about": "The states of the groups we want to list. If empty or null, all groups are returned with their state.", "fields": [ + { "name": "Name", "type": "string", "versions": "4+", "about": "The name of the group state" } ]} ] } diff --git a/clients/src/main/resources/common/message/ListGroupsResponse.json b/clients/src/main/resources/common/message/ListGroupsResponse.json index dcc96b6fb6f4f..c389d87e8a93c 100644 --- a/clients/src/main/resources/common/message/ListGroupsResponse.json +++ b/clients/src/main/resources/common/message/ListGroupsResponse.json @@ -23,7 +23,7 @@ // // Version 3 is the first flexible version. // - // Version 4 adds the GroupState flexible field (KIP-518). + // Version 4 adds the GroupState field (KIP-518). "validVersions": "0-4", "flexibleVersions": "3+", "fields": [ @@ -37,8 +37,8 @@ "about": "The group ID." }, { "name": "ProtocolType", "type": "string", "versions": "0+", "about": "The group protocol type." }, - { "name": "GroupState", "type": "string", "versions": "4+", "tag": 0, "taggedVersions": "4+", "ignorable": true, - "nullableVersions": "4+", "default": "null", "about": "The group state string." } + { "name": "GroupState", "type": "string", "versions": "4+", "nullableVersions": "0+", "ignorable": true, "default": "null", + "about": "The group state name." } ]} ] } diff --git a/clients/src/test/java/org/apache/kafka/clients/MockClient.java b/clients/src/test/java/org/apache/kafka/clients/MockClient.java index eaf5dcb4f8fd1..8948e9296aedb 100644 --- a/clients/src/test/java/org/apache/kafka/clients/MockClient.java +++ b/clients/src/test/java/org/apache/kafka/clients/MockClient.java @@ -215,15 +215,22 @@ public void send(ClientRequest request, long now) { AbstractRequest.Builder builder = request.requestBuilder(); short version = nodeApiVersions.latestUsableVersion(request.apiKey(), builder.oldestAllowedVersion(), builder.latestAllowedVersion()); - AbstractRequest abstractRequest = request.requestBuilder().build(version); - if (!futureResp.requestMatcher.matches(abstractRequest)) - throw new IllegalStateException("Request matcher did not match next-in-line request " + abstractRequest + " with prepared response " + futureResp.responseBody); UnsupportedVersionException unsupportedVersionException = null; if (futureResp.isUnsupportedRequest) - unsupportedVersionException = new UnsupportedVersionException("Api " + - request.apiKey() + " with version " + version); - + unsupportedVersionException = new UnsupportedVersionException( + "Api " + request.apiKey() + " with version " + version); + try { + AbstractRequest abstractRequest = request.requestBuilder().build(version); + if (!futureResp.requestMatcher.matches(abstractRequest)) + throw new IllegalStateException("Request matcher did not match next-in-line request " + + abstractRequest + " with prepared response " + futureResp.responseBody); + + } catch (UnsupportedVersionException uve) { + if (unsupportedVersionException == null) { + throw uve; + } + } ClientResponse resp = new ClientResponse(request.makeHeader(version), request.callback(), request.destination(), request.createdTimeMs(), time.milliseconds(), futureResp.disconnected, unsupportedVersionException, null, futureResp.responseBody); diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index 51686ffc41b62..742d1d84b6278 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.clients.admin; +import org.apache.kafka.clients.ApiVersion; import org.apache.kafka.clients.ClientDnsLookup; import org.apache.kafka.clients.ClientUtils; import org.apache.kafka.clients.MockClient; @@ -59,6 +60,7 @@ import org.apache.kafka.common.errors.UnknownMemberIdException; import org.apache.kafka.common.errors.UnknownServerException; import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; +import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData; import org.apache.kafka.common.message.CreatePartitionsResponseData; import org.apache.kafka.common.message.CreatePartitionsResponseData.CreatePartitionsTopicResult; @@ -91,6 +93,7 @@ import org.apache.kafka.common.message.OffsetDeleteResponseData.OffsetDeleteResponsePartitionCollection; import org.apache.kafka.common.message.OffsetDeleteResponseData.OffsetDeleteResponseTopic; import org.apache.kafka.common.message.OffsetDeleteResponseData.OffsetDeleteResponseTopicCollection; +import org.apache.kafka.common.protocol.ApiKeys; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.quota.ClientQuotaAlteration; import org.apache.kafka.common.quota.ClientQuotaEntity; @@ -116,6 +119,7 @@ import org.apache.kafka.common.requests.FindCoordinatorResponse; import org.apache.kafka.common.requests.IncrementalAlterConfigsResponse; import org.apache.kafka.common.requests.LeaveGroupResponse; +import org.apache.kafka.common.requests.ListGroupsRequest; import org.apache.kafka.common.requests.ListGroupsResponse; import org.apache.kafka.common.requests.ListOffsetResponse; import org.apache.kafka.common.requests.ListOffsetResponse.PartitionData; @@ -1216,10 +1220,12 @@ public void testListConsumerGroups() throws Exception { .setGroups(Arrays.asList( new ListGroupsResponseData.ListedGroup() .setGroupId("group-1") - .setProtocolType(ConsumerProtocol.PROTOCOL_TYPE), + .setProtocolType(ConsumerProtocol.PROTOCOL_TYPE) + .setGroupState("Stable"), new ListGroupsResponseData.ListedGroup() .setGroupId("group-connect-1") .setProtocolType("connector") + .setGroupState("Stable") ))), env.cluster().nodeById(0)); @@ -1245,10 +1251,12 @@ public void testListConsumerGroups() throws Exception { .setGroups(Arrays.asList( new ListGroupsResponseData.ListedGroup() .setGroupId("group-2") - .setProtocolType(ConsumerProtocol.PROTOCOL_TYPE), + .setProtocolType(ConsumerProtocol.PROTOCOL_TYPE) + .setGroupState("Stable"), new ListGroupsResponseData.ListedGroup() .setGroupId("group-connect-2") .setProtocolType("connector") + .setGroupState("Stable") ))), env.cluster().nodeById(1)); @@ -1259,10 +1267,12 @@ public void testListConsumerGroups() throws Exception { .setGroups(Arrays.asList( new ListGroupsResponseData.ListedGroup() .setGroupId("group-3") - .setProtocolType(ConsumerProtocol.PROTOCOL_TYPE), + .setProtocolType(ConsumerProtocol.PROTOCOL_TYPE) + .setGroupState("Stable"), new ListGroupsResponseData.ListedGroup() .setGroupId("group-connect-3") .setProtocolType("connector") + .setGroupState("Stable") ))), env.cluster().nodeById(2)); @@ -1283,7 +1293,7 @@ public void testListConsumerGroups() throws Exception { Set groupIds = new HashSet<>(); for (ConsumerGroupListing listing : listings) { groupIds.add(listing.groupId()); - assertFalse(listing.state().isPresent()); + assertTrue(listing.state().isPresent()); } assertEquals(Utils.mkSet("group-1", "group-2", "group-3"), groupIds); @@ -1347,6 +1357,46 @@ public void testListConsumerGroupsWithStates() throws Exception { } } + @Test + public void testListConsumerGroupsWithStatesOlderBrokerVersion() throws Exception { + ApiVersion listGroupV3 = new ApiVersion(ApiKeys.LIST_GROUPS.id, (short) 0, (short) 3); + try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(mockCluster(1, 0))) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create(Collections.singletonList(listGroupV3))); + + env.kafkaClient().prepareResponse(prepareMetadataResponse(env.cluster(), Errors.NONE)); + + // Check we can list groups with older broker if we don't specify states + env.kafkaClient().prepareResponseFrom( + new ListGroupsResponse(new ListGroupsResponseData() + .setErrorCode(Errors.NONE.code()) + .setGroups(Collections.singletonList( + new ListGroupsResponseData.ListedGroup() + .setGroupId("group-1") + .setProtocolType(ConsumerProtocol.PROTOCOL_TYPE)))), + env.cluster().nodeById(0)); + ListConsumerGroupsOptions options = new ListConsumerGroupsOptions().inAnyState(); + ListConsumerGroupsResult result = env.adminClient().listConsumerGroups(options); + Collection listing = result.all().get(); + assertEquals(1, listing.size()); + List expected = Collections.singletonList(new ConsumerGroupListing("group-1", false, Optional.empty())); + assertEquals(expected, listing); + + // But we cannot set a state filter with older broker + env.kafkaClient().prepareResponse(prepareMetadataResponse(env.cluster(), Errors.NONE)); + env.kafkaClient().prepareUnsupportedVersionResponse( + body -> body instanceof ListGroupsRequest); + + options = new ListConsumerGroupsOptions().inStates(Collections.singleton(ConsumerGroupState.STABLE)); + result = env.adminClient().listConsumerGroups(options); + try { + result.all().get(); + fail("Should have thrown"); + } catch (ExecutionException ee) { + assertTrue(ee.getCause() instanceof UnsupportedVersionException); + } + } + } + @Test public void testOffsetCommitNumRetries() throws Exception { final Cluster cluster = mockCluster(3, 0); @@ -2314,8 +2364,6 @@ public void testRemoveMembersFromGroupRetryBackoff() throws Exception { AtomicLong firstAttemptTime = new AtomicLong(0); AtomicLong secondAttemptTime = new AtomicLong(0); - final TopicPartition tp1 = new TopicPartition("foo", 0); - mockClient.prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); env.kafkaClient().prepareResponse(body -> { diff --git a/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java b/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java index a591dbb07c8b8..c712c185d54ea 100644 --- a/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java +++ b/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java @@ -230,11 +230,13 @@ public void testSerialization() throws Exception { checkRequest(createLeaveGroupRequest(), true); checkErrorResponse(createLeaveGroupRequest(), new UnknownServerException(), true); checkResponse(createLeaveGroupResponse(), 0, true); + for (short v = ApiKeys.LIST_GROUPS.oldestVersion(); v <= ApiKeys.LIST_GROUPS.latestVersion(); v++) { checkRequest(createListGroupsRequest(v), false); checkErrorResponse(createListGroupsRequest(v), new UnknownServerException(), true); + checkResponse(createListGroupsResponse(v), v, true); } - checkResponse(createListGroupsResponse(), 0, true); + checkRequest(createDescribeGroupRequest(), true); checkErrorResponse(createDescribeGroupRequest(), new UnknownServerException(), true); checkResponse(createDescribeGroupResponse(), 0, true); @@ -816,7 +818,7 @@ public void testValidApiVersionsRequest() { @Test(expected = UnsupportedVersionException.class) public void testListGroupRequestV3FailsWithStates() { ListGroupsRequestData data = new ListGroupsRequestData() - .setStates(asList(ConsumerGroupState.STABLE.name())); + .setStatesFilter(asList(ConsumerGroupState.STABLE.name())); new ListGroupsRequest.Builder(data).build((short) 3); } @@ -1079,19 +1081,20 @@ private SyncGroupResponse createSyncGroupResponse(int version) { private ListGroupsRequest createListGroupsRequest(short version) { ListGroupsRequestData data = new ListGroupsRequestData(); if (version >= 4) - data.setStates(Arrays.asList("Stable")); + data.setStatesFilter(Arrays.asList("Stable")); return new ListGroupsRequest.Builder(data).build(version); } - private ListGroupsResponse createListGroupsResponse() { - return new ListGroupsResponse( - new ListGroupsResponseData() - .setErrorCode(Errors.NONE.code()) - .setGroups(Collections.singletonList( - new ListGroupsResponseData.ListedGroup() - .setGroupId("test-group") - .setProtocolType("consumer") - ))); + private ListGroupsResponse createListGroupsResponse(int version) { + ListGroupsResponseData.ListedGroup group = new ListGroupsResponseData.ListedGroup() + .setGroupId("test-group") + .setProtocolType("consumer"); + if (version >= 4) + group.setGroupState("Stable"); + ListGroupsResponseData data = new ListGroupsResponseData() + .setErrorCode(Errors.NONE.code()) + .setGroups(Collections.singletonList(group)); + return new ListGroupsResponse(data); } private DescribeGroupsRequest createDescribeGroupRequest() { diff --git a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala index 13cbde94c614f..a2571c07b6d5c 100755 --- a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala +++ b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala @@ -98,7 +98,7 @@ object ConsumerGroupCommand extends Logging { } def consumerGroupStatesFromString(input: String): List[ConsumerGroupState] = { - val parsedStates = input.split(',').map(s => ConsumerGroupState.parse(s.trim.toLowerCase.capitalize)).toSet.toList + val parsedStates = input.split(',').map(s => ConsumerGroupState.parse(s.trim)).toSet.toList if (parsedStates.contains(ConsumerGroupState.UNKNOWN)) { val validStates = ConsumerGroupState.values().filter(_ != ConsumerGroupState.UNKNOWN) throw new IllegalArgumentException(s"Invalid state list '$input'. Valid states are: ${validStates.mkString(", ")}") diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 3f9660a477757..a4ba82d6a893e 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -1399,7 +1399,11 @@ class KafkaApis(val requestChannel: RequestChannel, def handleListGroupsRequest(request: RequestChannel.Request): Unit = { val listGroupsRequest = request.body[ListGroupsRequest] - val states = listGroupsRequest.data.states.asScala.toList + val states = if (listGroupsRequest.data.statesFilter == null) + // Handle a null array the same as empty + List() + else + listGroupsRequest.data.statesFilter.asScala.toList def createResponse(throttleMs: Int, groups: List[GroupOverview], error: Errors): AbstractResponse = { new ListGroupsResponse(new ListGroupsResponseData() @@ -1408,8 +1412,7 @@ class KafkaApis(val requestChannel: RequestChannel, val listedGroup = new ListGroupsResponseData.ListedGroup() .setGroupId(group.groupId) .setProtocolType(group.protocolType) - if (!states.isEmpty) - listedGroup.setGroupState(group.state.toString) + .setGroupState(group.state.toString) listedGroup }.asJava) .setThrottleTimeMs(throttleMs) diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index be5773d5b3578..3276b7c1628de 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1063,7 +1063,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { TestUtils.waitUntilTrue(() => { val matching = client.listConsumerGroups.all.get.asScala.filter(group => group.groupId == testGroupId && - !group.state.isPresent) + group.state.get == ConsumerGroupState.STABLE) matching.size == 1 }, s"Expected to be able to list $testGroupId") diff --git a/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala b/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala index cd1a31b3191fa..1659f18b79a76 100644 --- a/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala +++ b/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala @@ -75,13 +75,19 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { TestUtils.waitUntilTrue(() => { foundListing = service.listConsumerGroupsWithState(List(ConsumerGroupState.STABLE)).toSet expectedListingStable == foundListing - }, s"Expected to show groups expectedListingStable, but found $foundListing") + }, s"Expected to show groups $expectedListingStable, but found $foundListing") } @Test def testConsumerGroupStatesFromString(): Unit = { - val result = ConsumerGroupCommand.consumerGroupStatesFromString("STABLE, stable, Stable, eMpTy") - assertEquals(List(ConsumerGroupState.STABLE, ConsumerGroupState.EMPTY), result) + var result = ConsumerGroupCommand.consumerGroupStatesFromString("Stable") + assertEquals(List(ConsumerGroupState.STABLE), result) + + result = ConsumerGroupCommand.consumerGroupStatesFromString("Stable, PreparingRebalance") + assertEquals(List(ConsumerGroupState.STABLE, ConsumerGroupState.PREPARING_REBALANCE), result) + + result = ConsumerGroupCommand.consumerGroupStatesFromString("Dead,CompletingRebalance,") + assertEquals(List(ConsumerGroupState.DEAD, ConsumerGroupState.COMPLETING_REBALANCE), result) try { ConsumerGroupCommand.consumerGroupStatesFromString("bad, wrong") @@ -90,13 +96,13 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { } try { - ConsumerGroupCommand.consumerGroupStatesFromString(" bad, ") + ConsumerGroupCommand.consumerGroupStatesFromString("stable") } catch { case e: IllegalArgumentException => //Expected } try { - ConsumerGroupCommand.consumerGroupStatesFromString(" bad, stable") + ConsumerGroupCommand.consumerGroupStatesFromString(" bad, Stable") } catch { case e: IllegalArgumentException => //Expected } @@ -127,7 +133,7 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { out.contains("STATE") && out.contains(simpleGroup) && out.contains(group) }, s"Expected to find $simpleGroup, $group and the header, but found $out") - cgcArgs = Array("--bootstrap-server", brokerList, "--list", "--state", "stable") + cgcArgs = Array("--bootstrap-server", brokerList, "--list", "--state", "Stable") TestUtils.waitUntilTrue(() => { out = TestUtils.grabConsoleOutput(ConsumerGroupCommand.main(cgcArgs)) out.contains("STATE") && out.contains(group) && out.contains("Stable") diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index 5a5d3edf68bd9..b173d41c3887c 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -1725,15 +1725,16 @@ class KafkaApisTest { EasyMock.verify(replicaManager) } + @Test def testListGroupsRequest(): Unit = { val overviews = List( GroupOverview("group1", "protocol1", "Stable"), - GroupOverview("goupp2", "qwerty", "Empty") + GroupOverview("group2", "qwerty", "Empty") ) val response = listGroupRequest(None, overviews) assertEquals(2, response.data.groups.size) - assertNull(response.data.groups.get(0).groupState) - assertNull(response.data.groups.get(1).groupState) + assertEquals("Stable", response.data.groups.get(0).groupState) + assertEquals("Empty", response.data.groups.get(1).groupState) } @Test @@ -1751,7 +1752,7 @@ class KafkaApisTest { val data = new ListGroupsRequestData() if (state.isDefined) - data.setStates(Collections.singletonList(state.get)) + data.setStatesFilter(Collections.singletonList(state.get)) val listGroupsRequest = new ListGroupsRequest.Builder(data).build() val requestChannelRequest = buildRequest(listGroupsRequest) From 499b15b40039927ec5b6a599024a082e001b456e Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Thu, 28 May 2020 13:38:23 +0100 Subject: [PATCH 11/14] Address some of the feedback --- .../kafka/clients/admin/ConsumerGroupListing.java | 11 +++++++++++ .../kafka/clients/admin/KafkaAdminClient.java | 2 +- .../clients/admin/ListConsumerGroupsOptions.java | 8 -------- .../kafka/clients/admin/KafkaAdminClientTest.java | 11 +++-------- .../kafka/coordinator/group/GroupCoordinator.scala | 2 +- core/src/main/scala/kafka/server/KafkaApis.scala | 4 ++-- .../coordinator/group/GroupCoordinatorTest.scala | 14 +++++++------- .../scala/unit/kafka/server/KafkaApisTest.scala | 2 +- 8 files changed, 26 insertions(+), 28 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java b/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java index 40d96129625a0..989fbfd9323b4 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java @@ -30,6 +30,16 @@ public class ConsumerGroupListing { private final boolean isSimpleConsumerGroup; private final Optional state; + /** + * Create an instance with the specified parameters. + * + * @param groupId Group Id + * @param isSimpleConsumerGroup If consumer group is simple or not. + */ + public ConsumerGroupListing(String groupId, boolean isSimpleConsumerGroup) { + this(groupId, isSimpleConsumerGroup, Optional.empty()); + } + /** * Create an instance with the specified parameters. * @@ -38,6 +48,7 @@ public class ConsumerGroupListing { * @param state The state of the consumer group */ public ConsumerGroupListing(String groupId, boolean isSimpleConsumerGroup, Optional state) { + Objects.requireNonNull(state); this.groupId = groupId; this.isSimpleConsumerGroup = isSimpleConsumerGroup; this.state = state; diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 471626cb1dabc..e74c04d34726d 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3052,7 +3052,7 @@ void handleResponse(AbstractResponse abstractResponse) { runnable.call(new Call("listConsumerGroups", deadline, new ConstantNodeIdProvider(node.id())) { @Override ListGroupsRequest.Builder createRequest(int timeoutMs) { - List states = (options.states().isEmpty()) + List states = options.states().isEmpty() ? null : options.states() .stream() diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java index 81fffe06c6714..e64add431a07d 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java @@ -43,14 +43,6 @@ public ListConsumerGroupsOptions inStates(Set states) { return this; } - /** - * All groups with their states will be returned by listConsumerGroups() - */ - public ListConsumerGroupsOptions inAnyState() { - this.states = Collections.emptySet(); - return this; - } - /** * Returns the list of States that are requested or empty if no states have been specified */ diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index 742d1d84b6278..ba578f349259d 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -1344,7 +1344,7 @@ public void testListConsumerGroupsWithStates() throws Exception { .setGroupState("Empty")))), env.cluster().nodeById(0)); - final ListConsumerGroupsOptions options = new ListConsumerGroupsOptions().inAnyState(); + final ListConsumerGroupsOptions options = new ListConsumerGroupsOptions(); final ListConsumerGroupsResult result = env.adminClient().listConsumerGroups(options); Collection listings = result.valid().get(); @@ -1374,7 +1374,7 @@ public void testListConsumerGroupsWithStatesOlderBrokerVersion() throws Exceptio .setGroupId("group-1") .setProtocolType(ConsumerProtocol.PROTOCOL_TYPE)))), env.cluster().nodeById(0)); - ListConsumerGroupsOptions options = new ListConsumerGroupsOptions().inAnyState(); + ListConsumerGroupsOptions options = new ListConsumerGroupsOptions(); ListConsumerGroupsResult result = env.adminClient().listConsumerGroups(options); Collection listing = result.all().get(); assertEquals(1, listing.size()); @@ -1388,12 +1388,7 @@ public void testListConsumerGroupsWithStatesOlderBrokerVersion() throws Exceptio options = new ListConsumerGroupsOptions().inStates(Collections.singleton(ConsumerGroupState.STABLE)); result = env.adminClient().listConsumerGroups(options); - try { - result.all().get(); - fail("Should have thrown"); - } catch (ExecutionException ee) { - assertTrue(ee.getCause() instanceof UnsupportedVersionException); - } + TestUtils.assertFutureThrows(result.all(), UnsupportedVersionException.class); } } diff --git a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala index 094d3df9a41da..76945b9df3a94 100644 --- a/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala +++ b/core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala @@ -797,7 +797,7 @@ class GroupCoordinator(val brokerId: Int, } } - def handleListGroups(states: List[String]): (Errors, List[GroupOverview]) = { + def handleListGroups(states: Set[String]): (Errors, List[GroupOverview]) = { if (!isActive.get) { (Errors.COORDINATOR_NOT_AVAILABLE, List[GroupOverview]()) } else { diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index a4ba82d6a893e..4231b91607e96 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -1401,9 +1401,9 @@ class KafkaApis(val requestChannel: RequestChannel, val listGroupsRequest = request.body[ListGroupsRequest] val states = if (listGroupsRequest.data.statesFilter == null) // Handle a null array the same as empty - List() + immutable.Set[String]() else - listGroupsRequest.data.statesFilter.asScala.toList + listGroupsRequest.data.statesFilter.asScala.toSet def createResponse(throttleMs: Int, groups: List[GroupOverview], error: Errors): AbstractResponse = { new ListGroupsResponse(new ListGroupsResponseData() diff --git a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala index 90eb934ef8bf5..9791cd61b8a84 100644 --- a/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala +++ b/core/src/test/scala/unit/kafka/coordinator/group/GroupCoordinatorTest.scala @@ -174,7 +174,7 @@ class GroupCoordinatorTest { assertEquals(Errors.COORDINATOR_LOAD_IN_PROGRESS, describeGroupError) // ListGroups - val (listGroupsError, _) = groupCoordinator.handleListGroups(List()) + val (listGroupsError, _) = groupCoordinator.handleListGroups(Set()) assertEquals(Errors.COORDINATOR_LOAD_IN_PROGRESS, listGroupsError) // DeleteGroups @@ -3208,7 +3208,7 @@ class GroupCoordinatorTest { val syncGroupResult = syncGroupLeader(groupId, generationId, assignedMemberId, Map(assignedMemberId -> Array[Byte]())) assertEquals(Errors.NONE, syncGroupResult.error) - val (error, groups) = groupCoordinator.handleListGroups(List()) + val (error, groups) = groupCoordinator.handleListGroups(Set()) assertEquals(Errors.NONE, error) assertEquals(1, groups.size) assertEquals(GroupOverview("groupId", "consumer", Stable.toString), groups.head) @@ -3220,7 +3220,7 @@ class GroupCoordinatorTest { val joinGroupResult = dynamicJoinGroup(groupId, memberId, protocolType, protocols) assertEquals(Errors.NONE, joinGroupResult.error) - val (error, groups) = groupCoordinator.handleListGroups(List()) + val (error, groups) = groupCoordinator.handleListGroups(Set()) assertEquals(Errors.NONE, error) assertEquals(1, groups.size) assertEquals(GroupOverview("groupId", "consumer", CompletingRebalance.toString), groups.head) @@ -3228,7 +3228,7 @@ class GroupCoordinatorTest { @Test def testListGroupsWithStates(): Unit = { - val allStates = List(PreparingRebalance, CompletingRebalance, Stable, Dead, Empty).map(s => s.toString) + val allStates = Set(PreparingRebalance, CompletingRebalance, Stable, Dead, Empty).map(s => s.toString) val memberId = JoinGroupRequest.UNKNOWN_MEMBER_ID // Member joins the group @@ -3238,7 +3238,7 @@ class GroupCoordinatorTest { assertEquals(Errors.NONE, joinGroupResult.error) // The group should be in CompletingRebalance - val (error, groups) = groupCoordinator.handleListGroups(List(CompletingRebalance.toString)) + val (error, groups) = groupCoordinator.handleListGroups(Set(CompletingRebalance.toString)) assertEquals(Errors.NONE, error) assertEquals(1, groups.size) val (error2, groups2) = groupCoordinator.handleListGroups(allStates.filterNot(s => s == CompletingRebalance.toString)) @@ -3251,7 +3251,7 @@ class GroupCoordinatorTest { assertEquals(Errors.NONE, syncGroupResult.error) // The group is now stable - val (error3, groups3) = groupCoordinator.handleListGroups(List(Stable.toString)) + val (error3, groups3) = groupCoordinator.handleListGroups(Set(Stable.toString)) assertEquals(Errors.NONE, error3) assertEquals(1, groups3.size) val (error4, groups4) = groupCoordinator.handleListGroups(allStates.filterNot(s => s == Stable.toString)) @@ -3264,7 +3264,7 @@ class GroupCoordinatorTest { verifyLeaveGroupResult(leaveGroupResults) // The group is now empty - val (error5, groups5) = groupCoordinator.handleListGroups(List(Empty.toString)) + val (error5, groups5) = groupCoordinator.handleListGroups(Set(Empty.toString)) assertEquals(Errors.NONE, error5) assertEquals(1, groups5.size) val (error6, groups6) = groupCoordinator.handleListGroups(allStates.filterNot(s => s == Empty.toString)) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index b173d41c3887c..ab803df11d22c 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -1757,7 +1757,7 @@ class KafkaApisTest { val requestChannelRequest = buildRequest(listGroupsRequest) val capturedResponse = expectNoThrottling() - val expectedStates = if (state.isDefined) List(state.get) else List() + val expectedStates: Set[String] = if (state.isDefined) Set(state.get) else Set() EasyMock.expect(groupCoordinator.handleListGroups(expectedStates)) .andReturn((Errors.NONE, overviews)) EasyMock.replay(groupCoordinator, clientRequestQuotaManager, requestChannel) From 56f839afcf1d4c956cb807f6d4cdf79cf8f88d55 Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Thu, 28 May 2020 19:40:02 +0100 Subject: [PATCH 12/14] Address some of the feedback --- .../kafka/clients/admin/KafkaAdminClient.java | 10 +++---- .../common/requests/ListGroupsRequest.java | 2 +- .../common/message/ListGroupsRequest.json | 7 ++--- .../org/apache/kafka/clients/MockClient.java | 9 ++---- .../kafka/admin/ConsumerGroupCommand.scala | 29 ++++++++++--------- .../api/PlaintextAdminIntegrationTest.scala | 2 +- .../admin/DescribeConsumerGroupTest.scala | 21 ++++++++++++++ .../kafka/admin/ListConsumerGroupTest.scala | 14 +++++---- 8 files changed, 57 insertions(+), 37 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index e74c04d34726d..72c498fb84c7a 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3052,12 +3052,10 @@ void handleResponse(AbstractResponse abstractResponse) { runnable.call(new Call("listConsumerGroups", deadline, new ConstantNodeIdProvider(node.id())) { @Override ListGroupsRequest.Builder createRequest(int timeoutMs) { - List states = options.states().isEmpty() - ? null - : options.states() - .stream() - .map(s -> s.toString()) - .collect(Collectors.toList()); + List states = options.states() + .stream() + .map(s -> s.toString()) + .collect(Collectors.toList()); return new ListGroupsRequest.Builder(new ListGroupsRequestData().setStatesFilter(states)); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java index 6ad4e9ce641e7..ce3938530b615 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/ListGroupsRequest.java @@ -46,7 +46,7 @@ public Builder(ListGroupsRequestData data) { @Override public ListGroupsRequest build(short version) { - if (data.statesFilter() != null && version < 4) { + if (!data.statesFilter().isEmpty() && version < 4) { throw new UnsupportedVersionException("The broker only supports ListGroups " + "v" + version + ", but we need v4 or newer to request groups by states."); } diff --git a/clients/src/main/resources/common/message/ListGroupsRequest.json b/clients/src/main/resources/common/message/ListGroupsRequest.json index 44e7423a0ec71..dbe6d9b6f123a 100644 --- a/clients/src/main/resources/common/message/ListGroupsRequest.json +++ b/clients/src/main/resources/common/message/ListGroupsRequest.json @@ -25,9 +25,8 @@ "validVersions": "0-4", "flexibleVersions": "3+", "fields": [ - { "name": "StatesFilter", "type": "[]string", "versions": "4+", "nullableVersions": "4+", "default": "null", - "about": "The states of the groups we want to list. If empty or null, all groups are returned with their state.", "fields": [ - { "name": "Name", "type": "string", "versions": "4+", "about": "The name of the group state" } - ]} + { "name": "StatesFilter", "type": "[]string", "versions": "4+", + "about": "The states of the groups we want to list. If empty all groups are returned with their state." + } ] } diff --git a/clients/src/test/java/org/apache/kafka/clients/MockClient.java b/clients/src/test/java/org/apache/kafka/clients/MockClient.java index 8948e9296aedb..6cfc4fd8d373f 100644 --- a/clients/src/test/java/org/apache/kafka/clients/MockClient.java +++ b/clients/src/test/java/org/apache/kafka/clients/MockClient.java @@ -217,19 +217,14 @@ public void send(ClientRequest request, long now) { builder.latestAllowedVersion()); UnsupportedVersionException unsupportedVersionException = null; - if (futureResp.isUnsupportedRequest) + if (futureResp.isUnsupportedRequest) { unsupportedVersionException = new UnsupportedVersionException( "Api " + request.apiKey() + " with version " + version); - try { + } else { AbstractRequest abstractRequest = request.requestBuilder().build(version); if (!futureResp.requestMatcher.matches(abstractRequest)) throw new IllegalStateException("Request matcher did not match next-in-line request " + abstractRequest + " with prepared response " + futureResp.responseBody); - - } catch (UnsupportedVersionException uve) { - if (unsupportedVersionException == null) { - throw uve; - } } ClientResponse resp = new ClientResponse(request.makeHeader(version), request.callback(), request.destination(), request.createdTimeMs(), time.milliseconds(), futureResp.disconnected, diff --git a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala index a2571c07b6d5c..6f4853130b9d1 100755 --- a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala +++ b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala @@ -97,8 +97,8 @@ object ConsumerGroupCommand extends Logging { } } - def consumerGroupStatesFromString(input: String): List[ConsumerGroupState] = { - val parsedStates = input.split(',').map(s => ConsumerGroupState.parse(s.trim)).toSet.toList + def consumerGroupStatesFromString(input: String): Set[ConsumerGroupState] = { + val parsedStates = input.split(',').map(s => ConsumerGroupState.parse(s.trim)).toSet if (parsedStates.contains(ConsumerGroupState.UNKNOWN)) { val validStates = ConsumerGroupState.values().filter(_ != ConsumerGroupState.UNKNOWN) throw new IllegalArgumentException(s"Invalid state list '$input'. Valid states are: ${validStates.mkString(", ")}") @@ -202,15 +202,15 @@ object ConsumerGroupCommand extends Logging { def listGroups(): Unit = { if (opts.options.has(opts.stateOpt)) { - val stateValue = opts.options.valueOf(opts.stateOpt) - val states = if (stateValue == null || stateValue.isEmpty) - allStates - else - consumerGroupStatesFromString(stateValue) - val listings = listConsumerGroupsWithState(states) - printGroupStates(listings.map(e => (e.groupId, e.state.get.toString))) - } else - listConsumerGroups().foreach(println(_)) + val stateValue = opts.options.valueOf(opts.stateOpt) + val states = if (stateValue == null || stateValue.isEmpty) + Set[ConsumerGroupState]() + else + consumerGroupStatesFromString(stateValue) + val listings = listConsumerGroupsWithState(states) + printGroupStates(listings.map(e => (e.groupId, e.state.get.toString))) + } else + listConsumerGroups().foreach(println(_)) } def listConsumerGroups(): List[String] = { @@ -219,9 +219,9 @@ object ConsumerGroupCommand extends Logging { listings.map(_.groupId).toList } - def listConsumerGroupsWithState(states: List[ConsumerGroupState]): List[ConsumerGroupListing] = { + def listConsumerGroupsWithState(states: Set[ConsumerGroupState]): List[ConsumerGroupListing] = { val listConsumerGroupsOptions = withTimeoutMs(new ListConsumerGroupsOptions()) - listConsumerGroupsOptions.inStates(new java.util.HashSet(states.asJava)) + listConsumerGroupsOptions.inStates(states.asJava) val result = adminClient.listConsumerGroups(listConsumerGroupsOptions) result.all.get.asScala.toList } @@ -1086,6 +1086,9 @@ object ConsumerGroupCommand extends Logging { CommandLineUtils.printUsageAndDie(parser, s"Option $describeOpt takes at most one of these options: ${mutuallyExclusiveOpts.mkString(", ")}") } + if (options.has(stateOpt) && options.valueOf(stateOpt) != null) + CommandLineUtils.printUsageAndDie(parser, + s"Option $describeOpt does not take a value for $stateOpt") } else { if (options.has(timeoutMsOpt)) debug(s"Option $timeoutMsOpt is applicable only when $describeOpt is used.") diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 3276b7c1628de..43a6fa7aed952 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1068,7 +1068,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { }, s"Expected to be able to list $testGroupId") TestUtils.waitUntilTrue(() => { - val options = new ListConsumerGroupsOptions().inAnyState + val options = new ListConsumerGroupsOptions() val matching = client.listConsumerGroups(options).all.get.asScala.filter(group => group.groupId == testGroupId && group.state.get == ConsumerGroupState.STABLE) diff --git a/core/src/test/scala/unit/kafka/admin/DescribeConsumerGroupTest.scala b/core/src/test/scala/unit/kafka/admin/DescribeConsumerGroupTest.scala index 79621a7d6abce..2c28ccf8733ad 100644 --- a/core/src/test/scala/unit/kafka/admin/DescribeConsumerGroupTest.scala +++ b/core/src/test/scala/unit/kafka/admin/DescribeConsumerGroupTest.scala @@ -72,6 +72,27 @@ class DescribeConsumerGroupTest extends ConsumerGroupCommandTest { assertTrue(exitMessage.get.contains("Option [describe] takes at most one of these options")) } + @Test + def testDescribeWithStateValue(): Unit = { + var exitStatus: Option[Int] = None + var exitMessage: Option[String] = None + Exit.setExitProcedure { (status, err) => + exitStatus = Some(status) + exitMessage = err + throw new RuntimeException + } + val cgcArgs = Array("--bootstrap-server", brokerList, "--describe", "--all-groups", "--state", "Stable") + try { + ConsumerGroupCommand.main(cgcArgs) + } catch { + case e: RuntimeException => //expected + } finally { + Exit.resetExitProcedure() + } + assertEquals(Some(1), exitStatus) + assertTrue(exitMessage.get.contains("Option [describe] does not take a value for [state]")) + } + @Test def testDescribeOffsetsOfNonExistingGroup(): Unit = { val group = "missing.group" diff --git a/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala b/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala index 1659f18b79a76..1a585b5fb8798 100644 --- a/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala +++ b/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala @@ -64,7 +64,7 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { var foundListing = Set.empty[ConsumerGroupListing] TestUtils.waitUntilTrue(() => { - foundListing = service.listConsumerGroupsWithState(ConsumerGroupState.values.toList).toSet + foundListing = service.listConsumerGroupsWithState(ConsumerGroupState.values.toSet).toSet expectedListing == foundListing }, s"Expected to show groups $expectedListing, but found $foundListing") @@ -73,7 +73,7 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { foundListing = Set.empty[ConsumerGroupListing] TestUtils.waitUntilTrue(() => { - foundListing = service.listConsumerGroupsWithState(List(ConsumerGroupState.STABLE)).toSet + foundListing = service.listConsumerGroupsWithState(Set(ConsumerGroupState.STABLE)).toSet expectedListingStable == foundListing }, s"Expected to show groups $expectedListingStable, but found $foundListing") } @@ -81,34 +81,38 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { @Test def testConsumerGroupStatesFromString(): Unit = { var result = ConsumerGroupCommand.consumerGroupStatesFromString("Stable") - assertEquals(List(ConsumerGroupState.STABLE), result) + assertEquals(Set(ConsumerGroupState.STABLE), result) result = ConsumerGroupCommand.consumerGroupStatesFromString("Stable, PreparingRebalance") - assertEquals(List(ConsumerGroupState.STABLE, ConsumerGroupState.PREPARING_REBALANCE), result) + assertEquals(Set(ConsumerGroupState.STABLE, ConsumerGroupState.PREPARING_REBALANCE), result) result = ConsumerGroupCommand.consumerGroupStatesFromString("Dead,CompletingRebalance,") - assertEquals(List(ConsumerGroupState.DEAD, ConsumerGroupState.COMPLETING_REBALANCE), result) + assertEquals(Set(ConsumerGroupState.DEAD, ConsumerGroupState.COMPLETING_REBALANCE), result) try { ConsumerGroupCommand.consumerGroupStatesFromString("bad, wrong") + fail("Should have thrown") } catch { case e: IllegalArgumentException => //Expected } try { ConsumerGroupCommand.consumerGroupStatesFromString("stable") + fail("Should have thrown") } catch { case e: IllegalArgumentException => //Expected } try { ConsumerGroupCommand.consumerGroupStatesFromString(" bad, Stable") + fail("Should have thrown") } catch { case e: IllegalArgumentException => //Expected } try { ConsumerGroupCommand.consumerGroupStatesFromString(" , ,") + fail("Should have thrown") } catch { case e: IllegalArgumentException => //Expected } From 5e18e5c03c29a33e9098418b98c2c67843624c5b Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Thu, 28 May 2020 21:53:15 +0100 Subject: [PATCH 13/14] Address comments --- .../java/org/apache/kafka/clients/admin/KafkaAdminClient.java | 2 +- .../src/main/resources/common/message/ListGroupsResponse.json | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 72c498fb84c7a..e992bfe8b8734 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -3063,7 +3063,7 @@ private void maybeAddConsumerGroup(ListGroupsResponseData.ListedGroup group) { String protocolType = group.protocolType(); if (protocolType.equals(ConsumerProtocol.PROTOCOL_TYPE) || protocolType.isEmpty()) { final String groupId = group.groupId(); - final Optional state = group.groupState() == null + final Optional state = group.groupState().equals("") ? Optional.empty() : Optional.of(ConsumerGroupState.parse(group.groupState())); final ConsumerGroupListing groupListing = new ConsumerGroupListing(groupId, protocolType.isEmpty(), state); diff --git a/clients/src/main/resources/common/message/ListGroupsResponse.json b/clients/src/main/resources/common/message/ListGroupsResponse.json index c389d87e8a93c..87561c2ab964d 100644 --- a/clients/src/main/resources/common/message/ListGroupsResponse.json +++ b/clients/src/main/resources/common/message/ListGroupsResponse.json @@ -37,7 +37,7 @@ "about": "The group ID." }, { "name": "ProtocolType", "type": "string", "versions": "0+", "about": "The group protocol type." }, - { "name": "GroupState", "type": "string", "versions": "4+", "nullableVersions": "0+", "ignorable": true, "default": "null", + { "name": "GroupState", "type": "string", "versions": "4+", "ignorable": true, "about": "The group state name." } ]} ] From acb64f32f6fc2353dcaec2be44d05c6cd1a77b2d Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Fri, 29 May 2020 10:30:35 +0100 Subject: [PATCH 14/14] Address last round of reviews --- .../clients/admin/ConsumerGroupListing.java | 3 +-- .../admin/ListConsumerGroupsOptions.java | 5 ++-- .../kafka/admin/ConsumerGroupCommand.scala | 2 -- .../api/PlaintextAdminIntegrationTest.scala | 2 +- .../kafka/admin/ListConsumerGroupTest.scala | 27 ++++++------------- 5 files changed, 13 insertions(+), 26 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java b/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java index 989fbfd9323b4..0abc3e01ca9de 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ConsumerGroupListing.java @@ -48,10 +48,9 @@ public ConsumerGroupListing(String groupId, boolean isSimpleConsumerGroup) { * @param state The state of the consumer group */ public ConsumerGroupListing(String groupId, boolean isSimpleConsumerGroup, Optional state) { - Objects.requireNonNull(state); this.groupId = groupId; this.isSimpleConsumerGroup = isSimpleConsumerGroup; - this.state = state; + this.state = Objects.requireNonNull(state); } /** diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java index e64add431a07d..9f1f38dd4a8e6 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ListConsumerGroupsOptions.java @@ -35,8 +35,9 @@ public class ListConsumerGroupsOptions extends AbstractOptions states = Collections.emptySet(); /** - * Only groups in these states will be returned by listConsumerGroups() - * If not set, all groups are returned with their states + * If states is set, only groups in these states will be returned by listConsumerGroups() + * Otherwise, all groups are returned. + * This operation is supported by brokers with version 2.6.0 or later. */ public ListConsumerGroupsOptions inStates(Set states) { this.states = (states == null) ? Collections.emptySet() : new HashSet<>(states); diff --git a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala index 6f4853130b9d1..61f2a57ff2ce5 100755 --- a/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala +++ b/core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala @@ -46,8 +46,6 @@ import joptsimple.OptionException object ConsumerGroupCommand extends Logging { - val allStates = ConsumerGroupState.values.toList - def main(args: Array[String]): Unit = { val opts = new ConsumerGroupCommandOptions(args) diff --git a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala index 43a6fa7aed952..20f27be8f1676 100644 --- a/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/PlaintextAdminIntegrationTest.scala @@ -1068,7 +1068,7 @@ class PlaintextAdminIntegrationTest extends BaseAdminIntegrationTest { }, s"Expected to be able to list $testGroupId") TestUtils.waitUntilTrue(() => { - val options = new ListConsumerGroupsOptions() + val options = new ListConsumerGroupsOptions().inStates(Set(ConsumerGroupState.STABLE).asJava) val matching = client.listConsumerGroups(options).all.get.asScala.filter(group => group.groupId == testGroupId && group.state.get == ConsumerGroupState.STABLE) diff --git a/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala b/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala index 1a585b5fb8798..3997c9ef53709 100644 --- a/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala +++ b/core/src/test/scala/unit/kafka/admin/ListConsumerGroupTest.scala @@ -23,6 +23,7 @@ import kafka.utils.TestUtils import org.apache.kafka.common.ConsumerGroupState import org.apache.kafka.clients.admin.ConsumerGroupListing import java.util.Optional +import org.scalatest.Assertions.assertThrows class ListConsumerGroupTest extends ConsumerGroupCommandTest { @@ -59,8 +60,8 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { val service = getConsumerGroupService(cgcArgs) val expectedListing = Set( - new ConsumerGroupListing(simpleGroup, true, Optional.of(ConsumerGroupState.EMPTY)), - new ConsumerGroupListing(group, false, Optional.of(ConsumerGroupState.STABLE))) + new ConsumerGroupListing(simpleGroup, true, Optional.of(ConsumerGroupState.EMPTY)), + new ConsumerGroupListing(group, false, Optional.of(ConsumerGroupState.STABLE))) var foundListing = Set.empty[ConsumerGroupListing] TestUtils.waitUntilTrue(() => { @@ -69,7 +70,7 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { }, s"Expected to show groups $expectedListing, but found $foundListing") val expectedListingStable = Set( - new ConsumerGroupListing(group, false, Optional.of(ConsumerGroupState.STABLE))) + new ConsumerGroupListing(group, false, Optional.of(ConsumerGroupState.STABLE))) foundListing = Set.empty[ConsumerGroupListing] TestUtils.waitUntilTrue(() => { @@ -89,32 +90,20 @@ class ListConsumerGroupTest extends ConsumerGroupCommandTest { result = ConsumerGroupCommand.consumerGroupStatesFromString("Dead,CompletingRebalance,") assertEquals(Set(ConsumerGroupState.DEAD, ConsumerGroupState.COMPLETING_REBALANCE), result) - try { + assertThrows[IllegalArgumentException] { ConsumerGroupCommand.consumerGroupStatesFromString("bad, wrong") - fail("Should have thrown") - } catch { - case e: IllegalArgumentException => //Expected } - try { + assertThrows[IllegalArgumentException] { ConsumerGroupCommand.consumerGroupStatesFromString("stable") - fail("Should have thrown") - } catch { - case e: IllegalArgumentException => //Expected } - try { + assertThrows[IllegalArgumentException] { ConsumerGroupCommand.consumerGroupStatesFromString(" bad, Stable") - fail("Should have thrown") - } catch { - case e: IllegalArgumentException => //Expected } - try { + assertThrows[IllegalArgumentException] { ConsumerGroupCommand.consumerGroupStatesFromString(" , ,") - fail("Should have thrown") - } catch { - case e: IllegalArgumentException => //Expected } }