DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
| Table of Contents |
|---|
Status
Current state: DraftUnder discussion
Discussion thread: here [Change the link from the KIP proposal email archive to your own email thread]
...
Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).
Motivation
This KIP will deal with administration of groups.
Public Interfaces
Briefly list any new interfaces that will be introduced as part of this proposal or any existing interfaces that will be removed or changed. The purpose of this section is to concisely call out the public contract that will come along with this feature.
A public interface is any change to the following:
Binary log format
The network protocol and api behavior
Any class in the public packages under clientsConfiguration, especially client configuration
org/apache/kafka/common/serialization
org/apache/kafka/common
org/apache/kafka/common/errors
org/apache/kafka/clients/producer
org/apache/kafka/clients/consumer (eventually, once stable)
Monitoring
Command line tools and arguments
- Anything else that will likely break existing users in some way when they upgrade
Proposed Changes
Describe the new thing you want to do in appropriate detail. This may be fairly extensive and have large subsections of its own. Or it may be a few sentences. Use judgement based on the scope of the change.
Compatibility, Deprecation, and Migration Plan
- What impact (if any) will there be on existing users?
- If we are changing behavior how will we phase out the older behavior?
- If we need special migration tools, describe them here.
- When will we remove the existing behavior?
Test Plan
Describe in few sentences how the KIP will be tested. We are mostly interested in system tests (since unit-tests are specific to implementation details). How will we know that the implementation works as expected? How will we know nothing broke?
Rejected Alternatives
The behavior of groups in Apache Kafka is more complicated and subtle than it first appears. To most users of Kafka, groups are synonymous with consumer groups. However, the "classic" consumer group protocol was extensible and there are several well-known extensions in use. For example, distributed workers in Kafka Connect also use groups as a coordination mechanism, and some applications such as schema registries have also built upon the consumer group protocol in interesting ways. These are all groups. Then, KIP-848 introduced the new consumer group protocol and modern consumer groups use this new protocol. KIP-932 introduces share groups, and there may well be additional types of group in the future.
All of these types of groups share a namespace for group IDs, but the manner in which you administer a group depends upon its type.
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.
This KIP tries to resolve some of these situations and make it easier to work out what’s going on with the groups on a cluster.
Public Interfaces
Client API changes
AdminClient
Add the following methods on the org.apache.kafka.client.admin.AdminClient interface.
| Method signature | Description |
|---|---|
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.
*/
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); |
ListGroupOptions
| Code Block |
|---|
package org.apache.kafka.client.admin;
import org.apache.kafka.common.GroupType;
/**
* Options for {@link Admin#listGroups(ListGroupsOptions)}.
*
* The API of this class is evolving, see {@link Admin} for details.
*/
@InterfaceStability.Evolving
public class ListGroupsOptions extends AbstractOptions<ListGroupsOptions> {
} |
ListGroupsResult
| Code Block |
|---|
package org.apache.kafka.clients.admin;
/**
* The result of the {@link Admin#listGroups(ListGroupsOptions)} call.
* <p>
* The API of this class is evolving, see {@link Admin} for details.
*/
@InterfaceStability.Evolving
public class ListGroupsResult {
/**
* Returns a future that yields either an exception, or the full set of group listings.
*/
public KafkaFuture<Collection<GroupListing>> all() {
}
/**
* Returns a future which yields just the valid listings.
*/
public KafkaFuture<Collection<GroupListing>> valid() {
}
/**
* Returns a future which yields just the errors which occurred.
*/
public KafkaFuture<Collection<Throwable>> errors() {
}
} |
GroupListing
| Code Block |
|---|
package org.apache.kafka.client.admin;
import org.apache.kafka.common.ShareGroupState;
/**
* A listing of a group in the cluster.
* <p>
* The API of this class is evolving, see {@link Admin} for details.
*/
@InterfaceStability.Evolving
public class GroupListing {
public GroupListing(String groupId, GroupType type, String protocolType);
/**
* The id of the group.
*/
public String groupId();
/**
* The group type.
*/
public GroupType type();
/**
* The protocol type.
*/
public String protocolType();
} |
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:
| Option | Description |
|---|---|
--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. |
--describe | Describe the details of the groups. |
--group-type <String: type> | Filters the groups based on group type. Valid types are: 'consumer' (consumer groups) and 'share' (share groups). |
--help | Print usage information. |
--list | List all groups. |
--version | Display Kafka version. |
Note that --group-type consumer actually matches all groups whose protocol is "consumer" , and --group-type share matches all groups whose type is Share. 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 |
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 |
To list all of the consumer groups:
| Code Block |
|---|
$ bin/kafka-groups.sh --bootstrap-server localhost:9092 --list --group-type consumer
old-consumer-group
new-consumer-group |
To describe all of the consumer groups:
| 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 |
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
A new option --create is added to this tool to create a share group. If the share group exists, the command succeeds. If the group does not exist, the share group is created. If the group exists but it's not a share group, the command fails.
| Code Block |
|---|
$ bin/kafka-share-groups.sh --bootstrap-server localhost:9092 --create --group NewShareGroup
Share group 'NewShareGroup' exists.
$ bin/kafka-share-groups.sh --bootstrap-server localhost:9092 --create --group ExistingShareGroup
Share group 'ExistingShareGroup' exists.
$ bin/kafka-share-groups.sh --bootstrap-server localhost:9092 --create --group ConsumerGroup
Error: Group 'ConsumerGroup' is not a share group. |
Under the covers, it uses the AlterShareGroupOffsets RPC with an empty Topics array.
Also, 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. |
Proposed Changes
This KIP introduces a command-line tool for displaying all of the groups and their types.
Next, it introduces a way to create a share group administratively. If you create Kafka resources such as topics as part of deploying an application, you can now create share groups in the same way.
Finally, in situations where command-line tools are used to administer a group of the wrong type, you’ll now be told the group type is wrong, rather than the group does not exist.
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.
The ListGroups RPC response returns three pieces of information for each group: group ID, type and protocol. For the common types of group, here is what they mean:
Type | Protocol | Meaning |
|---|---|---|
Classic |
| Consumer group with the "classic" consumer group protocol |
Consumer |
| Consumer group with the KIP-848 consumer group protocol |
Share |
| Share group |
Classic |
| Kafka Connect distributed worker cluster group |
Classic | Any other string | Other customization of "classic" consumer group protocol, such as a schema registry |
The new kafka-groups.sh tool makes all of this information available.
Creating groups
Groups are created dynamically on first use. For example, when you connect the first consumer in a consumer group, the group coordinator creates the group as a consumer group. This is convenient, but it does mean that you need to ensure that different types of group use distinct group IDs or the results will be unpredictable. This is because the group type depends upon how the group was created, whether it was a consumer group member, a distributed Kafka Connect cluster, or whatever.
Today, you can create a consumer group administratively before the first member joins by resetting the offsets, such as like this:
$ bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --reset-offsets --to-latest --group CG1 --topic T1 --execute
That’s a slightly contorted way to create a consumer group, but it does already exist and it is known. As a result, this KIP does not introduce a new way to create consumer groups.
For share groups, a new --create option is added to the kafka-share-groups.sh tool to create a share group. You can use --reset-offsets in a similar way as consumer groups, but if you want to create a share group without specifying any topics, use --create instead. Note that the starting offset of a share group comes from the group's group.share.auto.offset.reset configuration property as introduced in KIP-932. This defaults to "latest" .
$ bin/kafka-share-groups.sh --bootstrap-server localhost:9092 --create --group SG1
Compatibility, Deprecation, and Migration Plan
The only 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.
Test Plan
The feature will be thoroughly tested with unit and integration tests.
Rejected Alternatives
NoneIf there are alternative ways of accomplishing the same thing, what were they? The purpose of this section is to motivate why the design is the way it is and not some other way.