This is an automated email from the ASF dual-hosted git repository.
dajac pushed a commit to branch 4.0
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/4.0 by this push:
new 6b645f8c6bf KAFKA-20640: Consumer group should reject classic member
joins when migration policy is disabled (#22457)
6b645f8c6bf is described below
commit 6b645f8c6bfe22c6d9e25f42759f67cb84bc831d
Author: Dongnuo Lyu <[email protected]>
AuthorDate: Fri Jun 12 08:04:05 2026 -0400
KAFKA-20640: Consumer group should reject classic member joins when
migration policy is disabled (#22457)
When group.consumer.migration.policy=disabled, a new classic-protocol
member is still admitted into an existing consumer group. Because the
policy permits neither an online upgrade nor an online downgrade, the
group can never converge back to a single protocol — it stays
permanently mixed.
This is inconsistent with the other direction. With disabled, a new
consumer-protocol member joining a classic group is refused, because the
upgrade it would require is turned off. A new classic member joining a
consumer group should be refused for the same reason: the downgrade it
would imply is also turned off. Under disabled the intent is that no
migration happens and groups stay single-protocol, so a new classic
member joining a consumer group should be rejected, mirroring the
consumer-into-classic case that is already rejected today.
This applies to disabled only. Under upgrade, classic members already in
the group during an in-progress upgrade must still be able to rejoin
after a transient fence, and under downgrade, new classic members are
the mechanism of the downgrade rollout, so neither policy should reject
them.
This patch rejects the new or any rejoining classic member if it tries
to join a consumer group when the policy is disabled.
Reviewers: David Jacot <[email protected]>
---
.../server/ConsumerProtocolMigrationTest.scala | 56 ++++++++++++-----
.../group/ConsumerGroupMigrationPolicy.java | 2 +-
.../coordinator/group/GroupCoordinatorConfig.java | 5 +-
.../coordinator/group/GroupMetadataManager.java | 19 ++++++
.../group/GroupMetadataManagerTest.java | 73 +++++++++++++++++++++-
5 files changed, 137 insertions(+), 18 deletions(-)
diff --git
a/core/src/test/scala/unit/kafka/server/ConsumerProtocolMigrationTest.scala
b/core/src/test/scala/unit/kafka/server/ConsumerProtocolMigrationTest.scala
index 0007f327146..01e12f45509 100644
--- a/core/src/test/scala/unit/kafka/server/ConsumerProtocolMigrationTest.scala
+++ b/core/src/test/scala/unit/kafka/server/ConsumerProtocolMigrationTest.scala
@@ -435,7 +435,7 @@ class ConsumerProtocolMigrationTest(cluster:
ClusterInstance) extends GroupCoord
new ClusterConfigProperty(key =
GroupCoordinatorConfig.CONSUMER_GROUP_MIGRATION_POLICY_CONFIG, value =
"disabled")
)
)
- def testDowngradeWithDisabledMigrationPolicy(): Unit = {
+ def testClassicMemberJoinToConsumerGroupWithDisabledMigrationPolicy(): Unit
= {
// Creates the __consumer_offsets topics because it won't be created
automatically
// in this test because it does not use FindCoordinator API.
createOffsetsTopic()
@@ -447,24 +447,50 @@ class ConsumerProtocolMigrationTest(cluster:
ClusterInstance) extends GroupCoord
)
val groupId = "grp"
+ val instanceId = "instance-id"
- // Consumer member 1 joins the group.
- val (memberId1, _) = joinConsumerGroupWithNewProtocol(groupId,
Uuid.randomUuid.toString)
+ // A static member using the consumer protocol joins the group.
+ consumerGroupHeartbeat(
+ groupId = groupId,
+ memberId = Uuid.randomUuid.toString,
+ instanceId = instanceId,
+ rebalanceTimeoutMs = 5 * 60 * 1000,
+ subscribedTopicNames = List("foo"),
+ topicPartitions = List.empty,
+ expectedError = Errors.NONE
+ )
- // Classic member 2 joins the group.
- val joinGroupResponseData = sendJoinRequest(
- groupId = groupId
+ val rejectedResponse = new JoinGroupResponseData()
+ .setProtocolName(null)
+ .setErrorCode(Errors.INCONSISTENT_GROUP_PROTOCOL.code)
+
+ // A new dynamic member is rejected.
+ assertEquals(
+ rejectedResponse,
+ sendJoinRequest(
+ groupId = groupId,
+ metadata = metadata(List.empty)
+ )
)
- sendJoinRequest(
- groupId = groupId,
- memberId = joinGroupResponseData.memberId,
- metadata = metadata(List.empty)
+
+ // A new static member with a different instance id is rejected.
+ assertEquals(
+ rejectedResponse,
+ sendJoinRequest(
+ groupId = groupId,
+ groupInstanceId = "another-instance-id",
+ metadata = metadata(List.empty)
+ )
)
- // Try to downgrade the group by leaving member 1.
- leaveGroupWithNewProtocol(
- groupId = groupId,
- memberId = memberId1
+ // A static member reusing the existing instance id is also rejected.
+ assertEquals(
+ rejectedResponse,
+ sendJoinRequest(
+ groupId = groupId,
+ groupInstanceId = instanceId,
+ metadata = metadata(List.empty)
+ )
)
// The group is still a consumer group.
@@ -473,7 +499,7 @@ class ConsumerProtocolMigrationTest(cluster:
ClusterInstance) extends GroupCoord
new ListGroupsResponseData.ListedGroup()
.setGroupId(groupId)
.setProtocolType("consumer")
- .setGroupState(ConsumerGroupState.ASSIGNING.toString)
+ .setGroupState(ConsumerGroupState.STABLE.toString)
.setGroupType(Group.GroupType.CONSUMER.toString)
),
listGroups(
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/ConsumerGroupMigrationPolicy.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/ConsumerGroupMigrationPolicy.java
index 7f86c8019e1..cf583fd79d7 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/ConsumerGroupMigrationPolicy.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/ConsumerGroupMigrationPolicy.java
@@ -33,7 +33,7 @@ public enum ConsumerGroupMigrationPolicy {
/** Only downgrade is enabled.*/
DOWNGRADE("downgrade", false, true),
- /** Neither upgrade nor downgrade is enabled.*/
+ /** Neither upgrade nor downgrade is enabled; a classic member cannot join
or rejoin a consumer group.*/
DISABLED("disabled", false, false);
private final String name;
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorConfig.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorConfig.java
index fb7c5444230..22ca9e2761b 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorConfig.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorConfig.java
@@ -205,7 +205,10 @@ public class GroupCoordinatorConfig {
ConsumerGroupMigrationPolicy.BIDIRECTIONAL + ": both upgrade from
classic group to consumer group and downgrade from consumer group to classic
group are enabled, " +
ConsumerGroupMigrationPolicy.UPGRADE + ": only upgrade from classic
group to consumer group is enabled, " +
ConsumerGroupMigrationPolicy.DOWNGRADE + ": only downgrade from
consumer group to classic group is enabled, " +
- ConsumerGroupMigrationPolicy.DISABLED + ": neither upgrade nor
downgrade is enabled.";
+ ConsumerGroupMigrationPolicy.DISABLED + ": neither upgrade nor
downgrade is enabled. " +
+ "In this mode a member using the classic protocol is not allowed to
join or rejoin a non-empty consumer group; " +
+ "if a group is already mixed (for example because this policy was
changed), its classic-protocol members are " +
+ "rejected when they next attempt to rejoin.";
///
/// Share group configs
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
index 7089bb11287..b35bfaa742a 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
@@ -1098,6 +1098,23 @@ public class GroupMetadataManager {
}
}
+ /**
+ * Validates whether a classic member is allowed to join or rejoin the
consumer group. When the
+ * migration policy is disabled, no classic member may join or rejoin the
consumer group.
+ *
+ * @param consumerGroup The consumer group the classic member wants to
join.
+ * @throws InconsistentGroupProtocolException if the classic member cannot
join the consumer group.
+ */
+ private void throwIfClassicMemberCannotJoinConsumerGroup(ConsumerGroup
consumerGroup) {
+ if (config.consumerGroupMigrationPolicy() ==
ConsumerGroupMigrationPolicy.DISABLED) {
+ log.info("Cannot join the consumer group {} with the classic
protocol because the group migration is disabled.",
+ consumerGroup.groupId());
+ throw Errors.INCONSISTENT_GROUP_PROTOCOL.exception(
+ String.format("Cannot join the consumer group %s with the
classic protocol because the group migration is disabled.",
consumerGroup.groupId())
+ );
+ }
+ }
+
/**
* Creates a ConsumerGroup corresponding to the given classic group.
*
@@ -1840,6 +1857,8 @@ public class GroupMetadataManager {
JoinGroupRequestData request,
CompletableFuture<JoinGroupResponseData> responseFuture
) throws ApiException {
+ throwIfClassicMemberCannotJoinConsumerGroup(group);
+
final long currentTimeMs = time.milliseconds();
final List<CoordinatorRecord> records = new ArrayList<>();
final String groupId = request.groupId();
diff --git
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
index 3445d6d8282..4516cc9ff27 100644
---
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
+++
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
@@ -11832,6 +11832,77 @@ public class GroupMetadataManagerTest {
assertThrows(InconsistentGroupProtocolException.class, () ->
context.sendClassicGroupJoin(requestWithInvalidProtocolType));
}
+ @ParameterizedTest
+ @ValueSource(booleans = {true, false})
+ public void
testClassicMemberJoinToConsumerGroupWithDisabledMigrationPolicy(boolean
isStatic) {
+ String groupId = "group-id";
+ String classicMemberId = Uuid.randomUuid().toString();
+ String classicInstanceId = "classic-instance-id";
+ String consumerMemberId = Uuid.randomUuid().toString();
+ String consumerInstanceId = "consumer-instance-id";
+ String errorMessage = String.format(
+ "Cannot join the consumer group %s with the classic protocol
because the group migration is disabled.", groupId);
+
+ // The group already contains a classic member. For the static case it
also contains a
+ // consumer-protocol static member that a classic member could try to
replace.
+ ConsumerGroupBuilder groupBuilder = new ConsumerGroupBuilder(groupId,
10)
+ .withMember(new ConsumerGroupMember.Builder(classicMemberId)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(10)
+ .setInstanceId(isStatic ? classicInstanceId : null)
+ .setClassicMemberMetadata(
+ new
ConsumerGroupMemberMetadataValue.ClassicMemberMetadata()
+ .setSessionTimeoutMs(5000)
+
.setSupportedProtocols(ConsumerGroupMember.classicProtocolListFromJoinRequestProtocolCollection(
+
GroupMetadataManagerTestContext.toProtocols("range"))))
+ .build());
+ if (isStatic) {
+ groupBuilder.withMember(new
ConsumerGroupMember.Builder(consumerMemberId)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(10)
+ .setInstanceId(consumerInstanceId)
+ .build());
+ }
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_MIGRATION_POLICY_CONFIG,
ConsumerGroupMigrationPolicy.DISABLED.toString())
+ .withConsumerGroup(groupBuilder)
+ .build();
+
+ // A new classic member cannot join
+ JoinGroupRequestData newMemberRequest = new
GroupMetadataManagerTestContext.JoinGroupRequestBuilder()
+ .withGroupId(groupId)
+ .withMemberId(UNKNOWN_MEMBER_ID)
+ .withGroupInstanceId(isStatic ? "new-instance-id" : null)
+
.withProtocols(GroupMetadataManagerTestContext.toProtocols("range"))
+ .build();
+ assertEquals(errorMessage,
+ assertThrows(InconsistentGroupProtocolException.class, () ->
context.sendClassicGroupJoin(newMemberRequest, isStatic)).getMessage());
+
+ // The existing classic member cannot rejoin.
+ JoinGroupRequestData rejoinRequest = new
GroupMetadataManagerTestContext.JoinGroupRequestBuilder()
+ .withGroupId(groupId)
+ .withMemberId(isStatic ? UNKNOWN_MEMBER_ID : classicMemberId)
+ .withGroupInstanceId(isStatic ? classicInstanceId : null)
+
.withProtocols(GroupMetadataManagerTestContext.toProtocols("range"))
+ .build();
+ assertEquals(errorMessage,
+ assertThrows(InconsistentGroupProtocolException.class, () ->
context.sendClassicGroupJoin(rejoinRequest, isStatic)).getMessage());
+
+ if (isStatic) {
+ // A classic member cannot replace an existing consumer-protocol
static member.
+ JoinGroupRequestData replaceRequest = new
GroupMetadataManagerTestContext.JoinGroupRequestBuilder()
+ .withGroupId(groupId)
+ .withMemberId(UNKNOWN_MEMBER_ID)
+ .withGroupInstanceId(consumerInstanceId)
+
.withProtocols(GroupMetadataManagerTestContext.toProtocols("range"))
+ .build();
+ assertEquals(errorMessage,
+ assertThrows(InconsistentGroupProtocolException.class, () ->
context.sendClassicGroupJoin(replaceRequest, true)).getMessage());
+ }
+ }
+
@Test
public void testJoiningConsumerGroupWithNewDynamicMember() throws
Exception {
String groupId = "group-id";
@@ -12129,7 +12200,7 @@ public class GroupMetadataManagerTest {
String memberId = Uuid.randomUuid().toString();
String instanceId = "instance-id";
GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
-
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_MIGRATION_POLICY_CONFIG,
ConsumerGroupMigrationPolicy.DISABLED.toString())
+
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_MIGRATION_POLICY_CONFIG,
ConsumerGroupMigrationPolicy.UPGRADE.toString())
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_ASSIGNORS_CONFIG, List.of(new
NoOpPartitionAssignor()))
.withMetadataImage(new MetadataImageBuilder()
.addTopic(fooTopicId, fooTopicName, 2)