DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
| Table of Contents |
|---|
Status
Current state: Under Discussion Accepted
Discussion thread: here, vote thread
...
Add new fields memberEpoch and protocolupgraded. For classic group member, there is no member epoch, so using Optional for the new field. For upgraded field, the value is like following:
- Classic Group -> Optional.empty
- Consumer Group -> Optional.empty if unknown
- Consumer Group -> Optional.of(false) if classic
- Consumer Group -> Optional.of(true) if consumer
Deprecated old constructors.
| Code Block | ||||||
|---|---|---|---|---|---|---|
| ||||||
public class MemberDescription {
private final String memberId;
private final Optional<String> groupInstanceId;
private final String clientId;
private final String host;
private final MemberAssignment assignment;
private final Optional<MemberAssignment> targetAssignment;
private final Optional<Integer> memberEpoch;
private final StringOptional<Boolean> protocolupgraded;
// ...
/**
* The epoch of the group member.
*/
public Optional<Integer> memberEpoch() {
return memberEpoch;
}
/**
public MemberDescription(String memberId,
* The protocol of the groupOptional<String> member.groupInstanceId,
*/
public String protocol() {
return protocol;
}
} |
RPCs
ConsumerGroupDescribeRequest
Bump validVersions to "0-1".
| Code Block | ||
|---|---|---|
| ||
{ "apiKey": 69, "type": "request", "listeners": ["broker"], "name": "ConsumerGroupDescribeRequest", "validVersions": "0-1", <-- bump to 0-1 "flexibleVersions": "0+", "fields": [ { "name": "GroupIds", "type": "[]string", "versions": "0+", "entityType": "groupId", String clientId, String host, "about": "The ids of the groups to describe" }, MemberAssignment assignment, { "name": "IncludeAuthorizedOperations", "type": "bool", "versions": "0+", "about": "Whether to include authorized operations." } ] } |
ConsumerGroupDescribeResponse
Bump validVersions to "0-1" and add a new field "Protocol" to "Members".
| Code Block | ||
|---|---|---|
| ||
{ "apiKey": 69, "type": "response", "name": "ConsumerGroupDescribeResponse", "validVersions": "0-1", <-- bump to 0-1 "flexibleVersions": "0+", // Supported errors: // - GROUP_AUTHORIZATION_FAILED (version 0+) // - NOT_COORDINATOR (version 0+) // - COORDINATOR_NOT_AVAILABLE (version 0+) // - COORDINATOR_LOAD_IN_PROGRESS (version 0+) // - INVALID_REQUEST (version 0+) // - INVALID_GROUP_ID (version 0+) // - GROUP_ID_NOT_FOUND (version 0+) "fields": [ { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+",Optional<MemberAssignment> targetAssignment, Optional<Integer> memberEpoch, Optional<Boolean> upgraded ) { this.memberId = memberId == null ? "" : memberId; this.groupInstanceId = groupInstanceId; "about": "The durationthis.clientId in= millisecondsclientId for== whichnull the? request"" was: throttledclientId; due to a quota violation, or zero ifthis.host the= requesthost did== notnull violate any quota." }, { "name": "Groups", "type": "[]DescribedGroup", "versions": "0+", ? "" : host; this.assignment = "about": "Each described group.", assignment == null ? "fields": [ { "name": "ErrorCode", "type": "int16", "versions": "0+",new MemberAssignment(Collections.emptySet()) : assignment; this.targetAssignment "about": "The describe error, or 0 if there was no error." },= targetAssignment; this.memberEpoch = memberEpoch; { "name": "ErrorMessage", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null",this.upgraded = upgraded; } /** * @deprecated Since 4.0. "about": "The top-level error message, or null if there was no error." },Use {@link #MemberDescription(String, Optional, String, String, MemberAssignment, Optional, Optional, Optional)}. */ @Deprecated { "name": "GroupId", "type": "string", "versions": "0+", "entityType": "groupId", public MemberDescription(String memberId, "about": "The group ID string." }, Optional<String> groupInstanceId, { "name": "GroupState", "type": "string", "versions": "0+", "about": "The group stateString stringclientId, or the empty string." }, { "name": "GroupEpoch", "type": "int32", "versions": "0+", "about": "The group epoch." }, String host, { "name": "AssignmentEpoch", "type": "int32", "versions": "0+", "about": "The assignment epoch." }MemberAssignment assignment, { "name": "AssignorName", "type": "string", "versions": "0+", "about": "The selected assignor." },Optional<MemberAssignment> targetAssignment ) { { "name": "Members", "type": "[]Member", "versions": "0+", this( "about": "The members."memberId, "fields": [ groupInstanceId, { "name": "MemberId", "type": "string", "versions": "0+", clientId, host, "about": "The member ID." }assignment, { "name": "InstanceId", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null"targetAssignment, Optional.empty(), Optional.empty() "about": "The member instance ID." }, ); } /** * { "name": "RackId", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null",@deprecated Since 4.0. Use {@link #MemberDescription(String, Optional, String, String, MemberAssignment, Optional, Optional, Optional)}. */ @Deprecated "about": "The member rack ID." },public MemberDescription( String memberId, { "name": "MemberEpoch", "type": "int32", "versions": "0+" Optional<String> groupInstanceId, String clientId, "about": "The current member epoch." }String host, MemberAssignment assignment ) { "name": "ClientId", "type": "string", "versions": "0+", this( memberId, "about": "The client ID." } groupInstanceId, { "name": "ClientHost", "type": "string", "versions": "0+", clientId, host, "about": "The client host." }assignment, { "name": "SubscribedTopicNames", "type": "[]string", "versions": "0+", "entityType": "topicName",Optional.empty() ); } /** "about": "The subscribed topic names." }, * @deprecated Since 4.0. Use {@link #MemberDescription(String, Optional, String, String, MemberAssignment, Optional, Optional, Optional)}. { "name": "SubscribedTopicRegex", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null", */ @Deprecated public MemberDescription(String memberId, "about": "the subscribed topic regex otherwise or null of not provided."String }clientId, { "name": "Assignment", "type": "Assignment", "versions": "0+", String host, "about": "The current assignment." }, { "name": "TargetAssignment", "type": "Assignment", "versions": "0+", MemberAssignment assignment) { "about": "The target assignment." }, { "name": "Protocol", "type": "string", "versions": "1+", "ignorable": true,this(memberId, Optional.empty(), clientId, host, assignment); } // ... /** * The epoch of the group member. */ public Optional<Integer> memberEpoch() { return memberEpoch; } /** * The flag indicating whether a member "about": "The member protocol." } <-- new fieldis classic. */ public Optional<Boolean> upgraded() { ]}, return upgraded; { "name": "AuthorizedOperations", } } |
RPCs
ConsumerGroupDescribeRequest
Bump validVersions to "0-1".
| Code Block | ||
|---|---|---|
| ||
{ "apiKey": 69, "type": "int32request", "versionslisteners": ["0+broker"], "defaultname": "-2147483648ConsumerGroupDescribeRequest", "validVersions": "0-1", <-- bump to 0-1 "aboutflexibleVersions": "32-bit bitfield to represent authorized operations for this group." } ] } ]0+", "commonStructsfields": [ { "name": "GroupIds", "TopicPartitionstype": "[]string", "versions": "0+", "fieldsentityType": ["groupId", { "nameabout": "TopicId", "type": "uuid", "versions": "0+", "about": "The topic ID." The ids of the groups to describe" }, { "name": "TopicNameIncludeAuthorizedOperations", "type": "stringbool", "versions": "0+", "entityType": "topicName", "about": "The topic nameWhether to include authorized operations." }, ] } |
ConsumerGroupDescribeResponse
Bump validVersions to "0-1" and add a new field "MemberType" to "Members". The "MemberType" field is int8. It's -1 as default for unknown, 0 for classic member, and 1 for consumer member.
| Code Block | ||
|---|---|---|
| ||
{ "apiKey": 69, "type": "response", { "name": "Partitions", "type": "[]int32", "versions": "0+", "about": "The partitions." } ]}, { "name": "AssignmentConsumerGroupDescribeResponse", "versionsvalidVersions": "0+-1", "fields": [ { "name <-- bump to 0-1 "flexibleVersions": "TopicPartitions0+", "type": "[]TopicPartitions", "versions": "0+", "about": "The assigned topic-partitions to the member." } ]} // Supported errors: // - GROUP_AUTHORIZATION_FAILED (version 0+) // - NOT_COORDINATOR (version 0+) // - COORDINATOR_NOT_AVAILABLE (version 0+) // - COORDINATOR_LOAD_IN_PROGRESS (version 0+) // - INVALID_REQUEST (version 0+) // - INVALID_GROUP_ID (version 0+) // - GROUP_ID_NOT_FOUND (version 0+) "fields": [ { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+", "about": "The duration in milliseconds for which the request was throttled due to a quota violation, or zero if the request did not violate any quota." }, { "name": "Groups", "type": "[]DescribedGroup", "versions": "0+", "about": "Each described group.", "fields": [ { "name": "ErrorCode", "type": "int16", "versions": "0+", "about": "The describe error, or 0 if there was no error." }, { "name": "ErrorMessage", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null", "about": "The top-level error message, or null if there was no error." }, { "name": "GroupId", "type": "string", "versions": "0+", "entityType": "groupId", "about": "The group ID string." }, { "name": "GroupState", "type": "string", "versions": "0+", "about": "The group state string, or the empty string." }, { "name": "GroupEpoch", "type": "int32", "versions": "0+", "about": "The group epoch." }, { "name": "AssignmentEpoch", "type": "int32", "versions": "0+", "about": "The assignment epoch." }, { "name": "AssignorName", "type": "string", "versions": "0+", "about": "The selected assignor." }, { "name": "Members", "type": "[]Member", "versions": "0+", "about": "The members.", "fields": [ { "name": "MemberId", "type": "string", "versions": "0+", "about": "The member ID." }, { "name": "InstanceId", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null", "about": "The member instance ID." }, { "name": "RackId", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null", "about": "The member rack ID." }, { "name": "MemberEpoch", "type": "int32", "versions": "0+", "about": "The current member epoch." }, { "name": "ClientId", "type": "string", "versions": "0+", "about": "The client ID." }, { "name": "ClientHost", "type": "string", "versions": "0+", "about": "The client host." }, { "name": "SubscribedTopicNames", "type": "[]string", "versions": "0+", "entityType": "topicName", "about": "The subscribed topic names." }, { "name": "SubscribedTopicRegex", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null", "about": "the subscribed topic regex otherwise or null of not provided." }, { "name": "Assignment", "type": "Assignment", "versions": "0+", "about": "The current assignment." }, { "name": "TargetAssignment", "type": "Assignment", "versions": "0+", "about": "The target assignment." }, { "name": "MemberType", "type": "int8", "versions": "1+", "default": "-1", "ignorable": true, "about": "-1 for unknown. 0 for classic member. +1 for consumer member." } <-- new field ]}, { "name": "AuthorizedOperations", "type": "int32", "versions": "0+", "default": "-2147483648", "about": "32-bit bitfield to represent authorized operations for this group." } ] } ], "commonStructs": [ { "name": "TopicPartitions", "versions": "0+", "fields": [ { "name": "TopicId", "type": "uuid", "versions": "0+", "about": "The topic ID." }, { "name": "TopicName", "type": "string", "versions": "0+", "entityType": "topicName", "about": "The topic name." }, { "name": "Partitions", "type": "[]int32", "versions": "0+", "about": "The partitions." } ]}, { "name": "Assignment", "versions": "0+", "fields": [ { "name": "TopicPartitions", "type": "[]TopicPartitions", "versions": "0+", "about": "The assigned topic-partitions to the member." } ]} ] } |
Proposed Changes
kafka-consumer-groups.sh
...
| Code Block |
|---|
$ bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group my-group --state --verbose GROUP COORDINATOR (ID) ASSIGNMENT-STRATEGY STATE GROUP-EPOCH TARGET-ASSIGNMENT-EPOCH #MEMBERS my-group localhost:61316 (0) uniform Stable 13 13 12 |
--describe --members --verbose
Show the member level level information. Add CURRENT-EPOCH, TARGET-EPOCH, TARGET-ASSIGNMENT, and UPGRADED information. Add The CURRENT-EPOCH , TARGET-EPOCH, TARGET-ASSIGNMENT, and PROTOCOL information. Change ASSIGNMENT to CURRENT-ASSIGNMENTmeans member epoch. Change ASSIGNMENT to CURRENT-ASSIGNMENT. For both classic and consumer groups, the ASSIGNMENT value format will change from "(0,1), (0,1)" to "my_topic:0,1;new_topic:0,1". The original format only shows sorted partition numbers for each topic. The new format adds the topic name. The UPGRADED is showed when a group is in migration.
| Code Block |
|---|
$ bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group my-group --members --verbose GROUP CONSUMER-ID HOST CLIENT-ID #PARTITIONS CURRENT-EPOCH CURRENT-ASSIGNMENT TARGET-EPOCH TARGET-ASSIGNMENT PROTOCOL my-group T4tbGKsvT7CtsxVBH5J2QQ CONSUMER-ID HOST CLIENT-ID #PARTITIONS CURRENT-EPOCH CURRENT-ASSIGNMENT TARGET-EPOCH TARGET-ASSIGNMENT UPGRADED my-group T4tbGKsvT7CtsxVBH5J2QQ /127.0.0.1 consumer-member 3 3 my_topic:0,1;new_topic:0 3 my_topic:0,1;new_topic:0 true my-group 6mfoIHq4n3BT7n1HCdqihb /127.0.0.1 consumer-my-group-1 4.1 classic-member 1 2 - 13 my_topic:0,1;new_topic:0,1 1 my_topic:0,1;new_topic:0,1 consumerfalse |
Following table lists protocol UPGRADED meaning.
| UPGRADED | Meaning |
|---|---|
| "true | |
| Protocol | Meaning |
| "classic" | A member in classic group. A member with non-null classicMemberMetadata in consumer group. |
| "consumer" | A consumer member with null classicMemberMetadatain a consumer group. |
| "sharefalse" | A classic member in share a consumer group. |
--describe --offsets --verbose / --describe --verbose
...
| Code Block |
|---|
$ bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group my-group --offsets --verbose GROUP TOPIC PARTITION LEADER-EPOCH CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID my-group my_topic 0 1 10 10 100 my-group-3056e339-59e1-43d8-abe3-8a9e10c53937T4tbGKsvT7CtsxVBH5J2QQ /127.0.0.1 consumer-test.group.CLASSIC.--offsets-2member my-group my_topic 1 1 10 10 100 my-group-3056e339-59e1-43d8-abe3-8a9e10c53937T4tbGKsvT7CtsxVBH5J2QQ /127.0.0.1 consumer-test.group.CLASSIC.--offsets-2 consumer-member my-group new_topic 0 1 10 10 100 my-group-3056e339-59e1-43d8-abe3-8a9e10c53937T4tbGKsvT7CtsxVBH5J2QQ /127.0.0.1 consumer-test.group.CLASSIC.--offsets-2member my-group new_topic 1 1 100 10 10 my-group-3056e339-59e1-43d8-abe3-8a9e10c539376mfoIHq4n3BT7n1HCdqihb /127.0.0.1 consumer-test.group.CLASSIC.--offsets-2 |
...
classic-member |
kafka-share-groups.sh
Add a new "–verbose" option to ShareGroupCommandOptions.
--describe --state --verbose
Show the group level information. Add GROUP-EPOCH and TARGET- ASSIGNMENT-EPOCH.
| Code Block |
|---|
$ bin/kafka-share-groups.sh --bootstrap-server localhost:9092 --describe --group my-group --state --verbose GROUP COORDINATOR (ID) STATE GROUP-EPOCH TARGET-ASSIGNMENT-EPOCH #MEMBERS my-group localhost:61316 (0) Stable 1 1 1 |
--describe --members --verbose
Show the member level information. Add CURRENTMEMBER-EPOCH , TARGET-EPOCH, and TARGET-ASSIGNMENT information. Change ASSIGNMENT to CURRENT-ASSIGNMENTand ASSIGNMENT information. The ASSIGNMENT value format will change from "my_topic:0,my_topic:1,new_topic:0,new_topic:1" to "my_topic:0,1;new_topic:0,1". The original format shows topic name for each partition. The new format only shows topic name one time for all partitions under the same topic.
| Code Block |
|---|
$ bin/kafka-share-groups.sh --bootstrap-server localhost:9092 --describe --group my-group --members --verbose GROUP CONSUMER-ID HOST CLIENT-ID CURRENTMEMBER-EPOCH CURRENT-ASSIGNMENT TARGET-EPOCH TARGET-ASSIGNMENT my-group T4tbGKsvT7CtsxVBH5J2QQ /127.0.0.1 consumer-my-group-1 1 my_topic:0,1;new_topic:0,1 1 my_topic:0,1;new_topic:0,1 |
...