dajac commented on code in PR #22356:
URL: https://github.com/apache/kafka/pull/22356#discussion_r3407898066
##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java:
##########
@@ -14292,6 +14463,75 @@ public void
testJoiningConsumerGroupThrowsExceptionIfGroupOverMaxSize() {
assertEquals("The consumer group has reached its maximum capacity of 1
members.", ex.getMessage());
}
+ @Test
+ public void
testStaticMemberCanRejoinConsumerGroupWithClassicProtocolWhenGroupIsFull()
throws Exception {
+ String groupId = "group-id";
+ String oldMemberId = "old-member";
+ String instanceId = "instance-id";
+
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withConsumerGroup(new ConsumerGroupBuilder(groupId, 10)
+ .withMember(new ConsumerGroupMember.Builder(oldMemberId)
+ .setInstanceId(instanceId)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(LEAVE_GROUP_STATIC_MEMBER_EPOCH)
+ .setPreviousMemberEpoch(9)
+ .build())
+ .withAssignmentEpoch(10))
+ .withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_MAX_SIZE_CONFIG,
1)
+
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_MIGRATION_POLICY_CONFIG,
ConsumerGroupMigrationPolicy.DISABLED.toString())
Review Comment:
This one too.
##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java:
##########
@@ -14292,6 +14463,75 @@ public void
testJoiningConsumerGroupThrowsExceptionIfGroupOverMaxSize() {
assertEquals("The consumer group has reached its maximum capacity of 1
members.", ex.getMessage());
}
+ @Test
+ public void
testStaticMemberCanRejoinConsumerGroupWithClassicProtocolWhenGroupIsFull()
throws Exception {
+ String groupId = "group-id";
+ String oldMemberId = "old-member";
+ String instanceId = "instance-id";
+
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withConsumerGroup(new ConsumerGroupBuilder(groupId, 10)
+ .withMember(new ConsumerGroupMember.Builder(oldMemberId)
+ .setInstanceId(instanceId)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(LEAVE_GROUP_STATIC_MEMBER_EPOCH)
+ .setPreviousMemberEpoch(9)
+ .build())
+ .withAssignmentEpoch(10))
+ .withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_MAX_SIZE_CONFIG,
1)
+
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_MIGRATION_POLICY_CONFIG,
ConsumerGroupMigrationPolicy.DISABLED.toString())
+ .build();
+
+ JoinGroupRequestData request = new
GroupMetadataManagerTestContext.JoinGroupRequestBuilder()
+ .withGroupId(groupId)
+ .withMemberId(UNKNOWN_MEMBER_ID)
+ .withGroupInstanceId(instanceId)
+
.withProtocols(GroupMetadataManagerTestContext.toConsumerProtocol(List.of(),
List.of()))
+ .build();
+
+ GroupMetadataManagerTestContext.JoinResult joinResult =
context.sendClassicGroupJoin(request, true, true);
+ joinResult.appendFuture.complete(null);
+ assertTrue(joinResult.joinFuture.isDone());
+
+ JoinGroupResponseData response = joinResult.joinFuture.get();
+ assertEquals(Errors.NONE.code(), response.errorCode());
+ assertNotEquals(UNKNOWN_MEMBER_ID, response.memberId());
+ assertNotEquals(oldMemberId, response.memberId());
+ assertEquals(response.memberId(),
context.groupMetadataManager.consumerGroup(groupId).staticMemberId(instanceId));
+ }
+
+ @Test
+ public void
testNewStaticMemberClassicGroupJoinThrowsGroupMaxSizeReachedExceptionWhenConsumerGroupIsFull()
throws Exception {
+ String groupId = "group-id";
+ String oldMemberId = "old-member";
+ String instanceId = "instance-id";
+ String newInstanceId = "new-instance-id";
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withConsumerGroup(new ConsumerGroupBuilder(groupId, 10)
+ .withMember(new ConsumerGroupMember.Builder(oldMemberId)
+ .setInstanceId(instanceId)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(LEAVE_GROUP_STATIC_MEMBER_EPOCH)
+ .setPreviousMemberEpoch(9)
+ .build())
+ .withAssignmentEpoch(10))
+ .withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_MAX_SIZE_CONFIG,
1)
+
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_MIGRATION_POLICY_CONFIG,
ConsumerGroupMigrationPolicy.DISABLED.toString())
Review Comment:
You must remove this config now.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]