Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.
Comment: Ready for review

...

Here’s an example of the complexity. If you start up a distributed Kafka Connect worker using the default configuration, it creates a group called "connect-cluster" . This is a group, but it’s not a consumer group. You can’t see this group in the list of consumer groups with the kafka-consumer-groups.sh  tool, but if you try to describe a consumer group called "connect-cluster"  or even use this group ID with a consumer, you get an error.

Here's another example. In Apache Kafka 3.8, if If you use kafka-consumer-groups.sh --describe --group MYSHARE  where the group is a share group, the output is Error: Consumer group 'MYSHARE' does not exist . That's probably OK, but it's interesting to see how it gets there.

First, the admin client uses the ConsumerGroupDescribe RPC which responds with error code GROUP_ID_NOT_FOUND (69)  and an empty error messagegiving no indication that it found a group of the wrong type. Next, the admin client falls back to the pre-KIP-848 DescribeGroups RPC in case it's a classic consumer group.  This RPC succeeds and responds with error code NONE (0)  and returns returning the group description with a status of Dead . It looks like a dead consumer group. Finally, this dead group is translated into the error message Error: Consumer group 'MYSHARE' does not exist . The output seems kind of acceptable, but the tool actually thinks it's dealing with a dead consumer group. It would be better if the ConsumerGroupDescribe RPC failed in a straightforward way.

...

This KIP introduces a command-line tool called kafka-groups.sh  for displaying all of the groups and their types.

In situations where command-line tools are used to administer a group of the wrong type, you’ll now be told the error message now indicates when the group type is wrong, rather than saying the group does not exist. This is achieved using a new error code in the Kafka protocol and a slight change in behavior of the admin client.

Listing groups

The KIP introduces a new tool called kafka-groups.sh to show all of the groups in a cluster, their types and the protocols they use. This lets you see consumer groups, share groups, Kafka Connect cluster groups, and any other custom groups all together. It doesn't replace the specific tools for the different types of group, but it does shed light on what's actually going on for administrators. Note that this does not require any changes to the Kafka protocol. The information is already available, but not directly accessible by the administrator.

...

Describing groups using the admin client

The behavior of A new exception InconsistentGroupTypeException  is used by AdminClient.describeConsumerGroups(Collection<String>) seems a little unusual. You can describe a collection of group IDs, some of which might exist and others might not. Here's how the response is built:

  1. If the group is a consumer group and the client is authorized to describe the group and there was no error, the group information is returned, along with the authorized operations if requested.
  2. If the group is not a consumer group (either does not exist or wrong type) and the client is authorized to describe the group and there was no error, the group information for a dead group is returned, along with the authorized operations if requested.
  3. If the client is not authorized to the describe the group, the group information error code is set to GROUP_AUTHORIZATION_FAILED . 
  4. If there was an error describing the group, the group information error code is set.

In cases (1) and (2), the admin client considers the operation a success, and this means the KafkaFuture for this group completes successfully. In cases (3) and (4), the admin client considers the operation unsuccessful, and this means the KafkaFuture for this group completes exceptionally. Case (2) is the tricky one because you can't readily tell what the dead group means. This is why the admin tools use output like Error: Consumer group 'MYSHARE" does not exist, even when the group ID is recognised and it's just the wrong type.

This gives a problem though because even though it's not difficult to determine which groups have the incorrect type, it would take a breaking change to the admin client to discover this situation if using the admin client.

As a result, this KIP introduces a new option on DescribeConsumerGroupOptions  called validateGroupType  which changes the behavior in the case where a group ID is the wrong group type. For case (2) above, the new INCONSISTENT_GROUP_TYPE  error in the ConsumerGroupDescribe response is translated into an InconsistentGroupTypeException containing the error message from the RPC response, and that can be used by the admin tools.

Public Interfaces

Client API changes

AdminClient

Add the following methods on the org.apache.kafka.client.admin.AdminClient  interface.

...

and AdminClient.describeShareGroups(Collection<String>)  to indicate that the group being described has the wrong type. This small change enables the error messages from the command line kafka-consumer-groups.sh  and kafka-share-groups.sh  to indicate when the group ID is found, but it has the wrong type.

Public Interfaces

Client API changes

AdminClient

Add the following methods on the org.apache.kafka.client.admin.AdminClient  interface.

Method signatureDescription
ListGroupsResult listGroups() List the groups available in the cluster.
ListGroupsResult listGroups(ListGroupsOptions options) List the groups available in the cluster.

Here are the method signatures:

Code Block
   /**
    * List the groups available in the cluster with the default options.
    *
    * <p>This is a convenience method for {@link #listGroups(ListGroupsOptions)} with default options.
    * See the overload for more details.
    *
    * @return The ListGroupsResult.
    *

Here are the method signatures:

Code Block
   /**
    * List the groups available in the cluster with the default options.
    *
    * <p>This is a convenience method for {@link #listGroups(ListGroupsOptions)} with default options.
    * See the overload for more details.
    *
    * @return The ListGroupsResult.
    */
   default ListGroupsResult listGroups() {
       return listGroups(new ListGroupsOptions());
   }
 
   /**
    * List the groups available in the cluster.
    *
    * @param options The options to use when listing the groups.
    * @return The ListGroupsResult.
    */
   ListGroupsResult listGroups(ListGroupsOptions options);

...

This class is modified to extend org.apache.kafka.clients.admin.GroupListing .

DescribeConsumerGroupsOptions

The following methods are added to this class.

Code Block
/**
 * Set to true if describing a group which is not a consumer group fails.
 */
public DescribeConsumerGroupsOptions validateGroupType(boolean validateGroupType);

/**
 * Set to true if describing a group which is not a consumer group fails.
 */
public boolean validateGroupType();

.GroupListing .

Exceptions

The following new exception is added to the org.apache.kafka.common.errors  package corresponding to the new error code in the Kafka protocol.

  • InconsistentGroupTypeException  - Indicates that the group exists but the group type is inconsistent with the operation.

Kafka protocol changes

This KIP adds the following error code to the Kafka protocol.

  • INCONSISTENT_GROUP_TYPE  (value TBD) - Indicates that the group exists but the group type is inconsistent with the operation.

This error code is used when the following RPCs are used with an existing group of the wrong type:

  • ConsumerGroupDescribe
  • ShareGroupDescribe

These RPCs are used by administrative tools and the new error code will help with the usability of the tools. A new version of ConsumerGroupDescribe with no schema change will be required to support the new error code. Because ShareGroupDescribe is still unstable, this error code can be added to the v0 RPC.

The remaining RPCs which work with groups, such as ListOffsets and TxnOffsetCommit, continue to fail with If validateGroupType  is set, when the ConsumerGroupDescribe RPC response contains the error code INCONSISTENT_GROUP_TYPE , the describe fails with InconsistentGroupTypeException . If it is not set, the error code is treated the same as GROUP_ID_NOT_FOUND  and the describe succeeds by returning a ConsumerGroupDescription  in Dead  state.

DescribeShareGroupsOptions

The following methods are added to this class.

Code Block
/**
 * Set to true if describing a group which is not a share group fails.
 */
public DescribeShareGroupsOptions validateGroupType(boolean validateGroupType);

/**
 * Set to true if describing a group which is not a share group fails.
 */
public boolean validateGroupType();

If validateGroupType  is set, when the ShareGroupDescribe RPC response contains the error code INCONSISTENT_GROUP_TYPE , the describe fails with InconsistentGroupTypeException . If it is not set, the error code is treated the same as GROUP_ID_NOT_FOUND  and the describe succeeds by returning a ShareGroupDescription  in Dead  state.

Exceptions

The following new exception is added to the org.apache.kafka.common.errors  package corresponding to the new error code in the Kafka protocol.

  • InconsistentGroupTypeException  - Indicates that the group exists but the group type is inconsistent with the operation.

Kafka protocol changes

This KIP adds the following error code to the Kafka protocol.

  • INCONSISTENT_GROUP_TYPE  (value TBD) - Indicates that the group exists but the group type is inconsistent with the operation.

This error code is used when the following RPCs are used with an existing group of the wrong type:

  • ConsumerGroupDescribe
  • ShareGroupDescribe

These RPCs are used by administrative tools and the new error code will help with the usability of the tools. Assuming that this KIP is delivered later than KIP-848, a new version of ConsumerGroupDescribe with no schema change will be required to support the new error code.

The remaining RPCs which work with groups, such as ListOffsets and TxnOffsetCommit, continue to fail with GROUP_ID_NOT_FOUND  if used against a group of the wrong type.

Command-line tools

kafka-groups.sh

A new tool called kafka-groups.sh  is introduced for listing and describing groups of any kind. It has the following options:

...

--bootstrap-server <String: server to connect to>

...

REQUIRED: The server(s) to connect to.

...

--command-config <String: command config property file>

...

Property file containing configs to be passed to Admin Client.

...

--consumer

...

Filters the groups based on group type and protocol in order to show consumer groups.

...

--describe

...

Describe the details of the groups.

...

--group-type <String: type>

...

Filters the groups based on group type. Valid types are: 'classic', 'consumer' and 'share'.

...

--help

...

Print usage information.

...

--list

...

List all groups.

...

--protocol <String: protocol>

...

Filters the groups based on protocol type.

...

--version

...

Display Kafka version.

if used against a group of the wrong type.

Command-line tools

kafka-groups.sh

A new tool called kafka-groups.sh  is introduced for listing and describing groups of any kind. It has the following options:

OptionDescription

--bootstrap-server <String: server to connect to>

REQUIRED: The server(s) to connect to.

--command-config <String: command config property file>

Property file containing configs to be passed to Admin Client.

--consumer

Filters the groups based on group type and protocol in order to show consumer groups.

--describe

Describe the details of the groups.

--group-type <String: type>

Filters the groups based on group type. Valid types are: 'classic', 'consumer' and 'share'.

--help

Print usage information.

--list

List all groups.

--protocol <String: protocol>

Filters the groups based on protocol type.

--version

Display Kafka version.

Note that --consumer  actually matches all groups whose type is Consumer , and groups whose type is Classic and protocol type is "consumer", and also "simple" consumer groups whose type is Classic and protocol type is "" . The filtering is done in the kafka-groups.sh  tool.

Here are some examples.

To list all of the groups:

Code Block
$ bin/kafka-groups.sh --bootstrap-server localhost:9092 --list
old-consumer-group
new-consumer-group
connect-cluster
share-group
schema-registry
simple-consumer-group

To describe all of the groups and their types:

Code Block
$ bin/kafka-groups.sh --bootstrap-server localhost:9092 --describe
GROUP                   TYPE          PROTOCOL
old-consumer-group      Classic       consumer
new-consumer-group      Consumer      consumer
connect-cluster         Classic       connect
share-group             Share         share
schema-registry         Classic       sr
simple-consumer-group   Classic

To list all of the consumer groups, silently merging together all of the different kinds of consumer group:

Code Block
$ bin/kafka-groups.sh --bootstrap-server localhost:9092 --list --consumer
old-consumer-group
new-consumer-group
simple-consumer-group

To describe all of the consumer groups:

Code Block
$ bin/kafka-groups.sh --bootstrap-server localhost:9092 --describe --consumer
GROUP                   TYPE          PROTOCOL
old-consumer-group      Classic       consumer
new-consumer-group      Consumer      consumer
simple-consumer-group   Classic

To describe all of the KIP-848 consumer groups

Note that --consumer  actually matches all groups whose type is Consumer , and groups whose type is Classic and protocol type is "consumer", and also "simple" consumer groups whose type is Classic and protocol type is "" . The filtering is done in the kafka-groups.sh  tool.

Here are some examples.

To list all of the groups:

Code Block
$ bin/kafka-groups.sh --bootstrap-server localhost:9092 --list
old-consumer-group
new-consumer-group
connect-cluster
share-group
schema-registry
simple-consumer-group

To describe all of the groups and their types:

Code Block
$ bin/kafka-groups.sh --bootstrap-server localhost:9092 --describe --group-type consumer
GROUP                   TYPE          PROTOCOL
old-consumer-group      Classic       consumer
new-consumer-group      Consumer      consumer
connect-cluster     

To describe all of the share groups:

Code Block
$ bin/kafka-groups.sh --bootstrap-server localhost:9092 --describe --group-type share
GROUP    Classic       connect
share-group        TYPE     Share         share
schema-registryPROTOCOL
share-group         Classic    Share   sr
simple-consumer-group   Classic

To list all of the consumer groups:

Code Block
$ bin/kafka-groups.sh --bootstrap-server localhost:9092 --list --consumer
old-consumer-group
new-consumer-group
simple-consumer-group
      share

kafka-consumer-groups.sh

For all operations which act on a single group, if that group exists but is not a consumer group, the command fails with a message indicating that the group type is incorrect, rather than the existing message that the group does not exist.

For example, if you try to describe a share group, the output will look like thisTo describe all of the consumer groups:

Code Block
$ bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --consumer
GROUP                   TYPE          PROTOCOL
old-consumer-group      Classic       consumer
new-consumer-group      Consumer      consumer
simple-consumer-group   Classic9092 --describe --group SG1
Error: Group 'SG1' is not a consumer group.

kafka-share-groups.sh

For all operations which act on a single group, if that group exists but is not a share group, the command fails with a message indicating that the group type is incorrect, rather than the existing message that the group does not exist.

For example, if you try to describe a consumer group, the output will look like thisTo describe all of the KIP-848 consumer groups:

Code Block
$ bin/kafka-share-groups.sh --bootstrap-server localhost:9092 --describe --group-type consumerCG1
GROUPError: Group 'CG1' is not a              TYPE          PROTOCOL
new-consumer-group      Consumer      consumer

To describe all of the share groups:

Code Block
$ bin/kafka-groups.sh --bootstrap-server localhost:9092 --describe --group-type share
GROUP                   TYPE          PROTOCOL
share-group             Share         share

kafka-consumer-groups.sh

For all operations which act on a single group, if that group exists but is not a consumer group, the command fails with a message indicating that the group type is incorrect, rather than the existing message that the group does not exist.

For example, if you try to describe a share group, the output will look like this:

Code Block
$ bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group SG1
Error: Group 'SG1' is not a consumer group.

kafka-share-groups.sh

For all operations which act on a single group, if that group exists but is not a share group, the command fails with a message indicating that the group type is incorrect, rather than the existing messages that the group does not exist.

For example, if you try to describe a consumer group, the output will look like this:

Code Block
$ bin/kafka-share-groups.sh --bootstrap-server localhost:9092 --describe --group CG1
Error: Group 'CG1' is not a share group.

Compatibility, Deprecation, and Migration Plan

A change to existing behavior is the error messages issued by the kafka-consumer-groups.sh  and kafka-share-groups.sh  tools when working with groups of the wrong type.

We also need to consider compatibility for the ListGroups RPC and the AdminClient.listGroups() method. ListGroups v5 introduces the TypesFilter  in the request and the GroupType  in the response. This KIP does not introduce a new version of the ListGroups RPC.

  • If the broker does not support ListGroups v5, setting a types filter using ListGroupOptions.withTypes() results in UnsupportedVersionException when the ListGroups request is serialized. (This is KIP-848 behavior.)
  • If the broker receives a group type in TypesFilter  that it does not support in a ListGroups v5 (or later) request, the filter is treated as unknown group type which thus matches no groups. (This is KIP-848 behavior.)
  • If the client receives a group type that it does not understand in a ListGroups v5 (or later) response, the group type is GroupType.UNKNOWN . This is because the group type in the ListGroups RPC response cannot be parsed into a value of the GroupType  enumeration, and the default value of UNKNOWN  is used in these situations.
share group.

Compatibility, Deprecation, and Migration Plan

ListGroups RPC

When writing this KIP, it seemed that perhaps the ListGroups RPC would need to be enhanced to create the AdminClient.listGroups()  method. This turned out not to be necessary.

ListGroups v5 introduces the TypesFilter  in the request and the GroupType  in the response. This KIP does not need to introduce a new version of the ListGroups RPC, because:

  • If the broker does not support ListGroups v5, setting a types filter using ListGroupOptions.withTypes() results in UnsupportedVersionException when the ListGroups request is serialized. (This is KIP-848 behavior.)
  • If the broker receives a group type in TypesFilter  that it does not support in a ListGroups v5 (or later) request, the filter is treated as unknown group type which thus matches no groups. (This is KIP-848 behavior.)
  • If the client receives a group type that it does not understand in a ListGroups v5 (or later) response, the group type is GroupType.UNKNOWN . This is because the group type in the ListGroups RPC response cannot be parsed into a value of the GroupType  enumeration, and the default value of UNKNOWN  is used in these situations.

The existing behavior suffices and there is no compatibility problem.

DescribeConsumerGroups

Prior to this KIP, the behavior of AdminClient.describeConsumerGroups(Collection<String>) seems a little unusual when used with groups which are not consumer groups. You can describe a collection of group IDs, some of which might exist and others might not. Here's how the response is built:

  1. If the group is a consumer group and the client is authorized to describe the group and there was no error, the group information is returned, along with the authorized operations if requested.
  2. If the group is not a consumer group (either does not exist or wrong type) and the client is authorized to describe the group and there was no error, the group information for a dead group is returned, along with the authorized operations if requested.
  3. If the client is not authorized to the describe the group, the group information error code is set to GROUP_AUTHORIZATION_FAILED . 
  4. If there was an error describing the group, the group information error code is set.

In cases (1) and (2), the admin client considers the operation a success, and this means the KafkaFuture for this group completes successfully. In cases (3) and (4), the admin client considers the operation unsuccessful, and this means the KafkaFuture for this group completes exceptionally. Case (2) is the tricky one because you can't readily tell what the dead group means. This is why the admin tools use output like Error: Consumer group 'MYSHARE" does not exist, even when the group ID is recognised and it's just the wrong type.

After this KIP, if you use AdminClient.describeConsumerGroup(Collection<String>)  to describe a group which is not a consumer group, case (2) above will result in the group information error code set to INCONSISTENT_GROUP_TYPE . In the admin client, this means the future for this group completes exceptionally with InconsistentGroupTypeException  rather than succeeding with a dead consumer groupAs a result, there is no compatibility problem with ListGroups as a result of this KIP.

Test Plan

The feature will be thoroughly tested with unit and integration tests.

Rejected Alternatives

NoneIt would be possible to preserve the current behavior of AdminClient.describeConsumerGroups(Collection<String>)  when used with a group which is not a consumer group and introduce an option to ask it to validate the group type rather than converting any indescribable group into a dead consumer group. This seems like an unnecessary complication with little benefit.