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)

Reply via email to