Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -1531,15 +1531,25 @@ private static boolean areOwnedTasksContainedInAssignedTasksWithEpochs(
* Checks whether the consumer group can accept a new member or not based on the
* max group size defined.
*
* @param group The consumer group.
* @param memberId The member id.
* @param group The consumer group.
* @param memberId The member id.
* @param instanceId The instance id.
*
* @throws GroupMaxSizeReachedException if the maximum capacity has been reached.
*/
private void throwIfConsumerGroupIsFull(
ConsumerGroup group,
String memberId
String memberId,
String instanceId
) throws GroupMaxSizeReachedException {
// If a static member already exists, we do not enforce the maximum group size check.
// An existing static member will fall into one of the following two cases,
// and neither affects the group size:
// 1. The member is replaced due to the static member rejoining.
// 2. 'UnreleasedInstanceIdException' is raised due to an epoch mismatch.
if (group.hasStaticMember(instanceId))
return;

// If the consumer group has reached its maximum capacity, the member is rejected if it is not
// already a member of the consumer group.
if (group.numMembers() >= config.consumerGroupMaxSize() && (memberId.isEmpty() || !group.hasMember(memberId))) {
Expand Down Expand Up @@ -2331,7 +2341,7 @@ private CoordinatorResult<ConsumerGroupHeartbeatResponseData, CoordinatorRecord>
// Get or create the consumer group.
boolean createIfNotExists = memberEpoch == 0;
final ConsumerGroup group = getOrMaybeCreateConsumerGroup(groupId, createIfNotExists, records);
throwIfConsumerGroupIsFull(group, memberId);
throwIfConsumerGroupIsFull(group, memberId, instanceId);

// Get or create the member.
if (memberId.isEmpty()) memberId = Uuid.randomUuid().toString();
Expand Down Expand Up @@ -2509,7 +2519,7 @@ private CoordinatorResult<Void, CoordinatorRecord> classicGroupJoinToConsumerGro
final boolean isUnknownMember = memberId.equals(UNKNOWN_MEMBER_ID);
if (isUnknownMember) memberId = Uuid.randomUuid().toString();

throwIfConsumerGroupIsFull(group, memberId);
throwIfConsumerGroupIsFull(group, memberId, instanceId);
throwIfClassicProtocolIsNotSupported(group, memberId, request.protocolType(), protocols);

if (JoinGroupRequest.requiresKnownMemberId(request, context.requestVersion())) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3701,6 +3701,177 @@ barTopicName, computeTopicHash(barTopicName, metadataImage)
.setTopicPartitions(List.of())));
}

@Test
public void testStaticMemberCanRejoinWhenConsumerGroupIsFull() {
String groupId = "fooup";
String instanceId = "instance-id";
String oldMemberId = "old-member-id";
String newMemberId = "new-member-id";
int groupMaxSize = 1;

Uuid fooTopicId = Uuid.randomUuid();
String fooTopicName = "foo";

CoordinatorMetadataImage metadataImage = new MetadataImageBuilder()
.addTopic(fooTopicId, fooTopicName, 1)
.buildCoordinatorMetadataImage();

MockPartitionAssignor assignor = new MockPartitionAssignor("range");
assignor.prepareGroupAssignment(new GroupAssignment(Map.of(
newMemberId, new MemberAssignmentImpl(mkAssignment(mkTopicAssignment(fooTopicId, 0)))
)));
GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder()
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_ASSIGNORS_CONFIG, List.of(assignor))
.withMetadataImage(metadataImage)
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_MAX_SIZE_CONFIG, groupMaxSize)
.withConsumerGroup(new ConsumerGroupBuilder(groupId, 10)
.withMember(new ConsumerGroupMember.Builder(oldMemberId)
.setState(MemberState.STABLE)
.setInstanceId(instanceId)
.setMemberEpoch(LEAVE_GROUP_STATIC_MEMBER_EPOCH)
.setPreviousMemberEpoch(9)
.setRebalanceTimeoutMs(5000)
.setSubscribedTopicNames(List.of("foo", "bar"))
.setServerAssignorName("range")
.setAssignedPartitions(mkAssignment(
mkTopicAssignment(fooTopicId, 0)))
.build())
.withAssignment(oldMemberId, mkAssignment(mkTopicAssignment(fooTopicId, 0)))
.withAssignmentEpoch(10)
.withMetadataHash(computeGroupHash(Map.of(fooTopicName, computeTopicHash(fooTopicName, metadataImage)))))
.build();

CoordinatorResult<ConsumerGroupHeartbeatResponseData, CoordinatorRecord> result = context.consumerGroupHeartbeat(
new ConsumerGroupHeartbeatRequestData()
.setGroupId(groupId)
.setInstanceId(instanceId)
.setMemberId(newMemberId)
.setMemberEpoch(0)
.setServerAssignor("range")
.setRebalanceTimeoutMs(5000)
.setSubscribedTopicNames(List.of("foo"))
.setTopicPartitions(List.of()));

assertResponseEquals(
new ConsumerGroupHeartbeatResponseData()
.setMemberId(newMemberId)
.setMemberEpoch(11)
.setHeartbeatIntervalMs(5000)
.setAssignment(new ConsumerGroupHeartbeatResponseData.Assignment()
.setTopicPartitions(List.of(
new ConsumerGroupHeartbeatResponseData.TopicPartitions()
.setTopicId(fooTopicId)
.setPartitions(List.of(0))
))),
result.response()
);
}

@Test
public void testStaticMemberRejoinWithUnreleasedInstanceIdFailsWhenConsumerGroupIsFull() {
String groupId = "fooup";
String instanceId = "instance-id";
String oldMemberId = "old-member-id";
String newMemberId = "new-member-id";
int groupMaxSize = 1;

Uuid fooTopicId = Uuid.randomUuid();
String fooTopicName = "foo";

CoordinatorMetadataImage metadataImage = new MetadataImageBuilder()
.addTopic(fooTopicId, fooTopicName, 1)
.buildCoordinatorMetadataImage();

MockPartitionAssignor assignor = new MockPartitionAssignor("range");
assignor.prepareGroupAssignment(new GroupAssignment(Map.of(
oldMemberId, new MemberAssignmentImpl(mkAssignment(mkTopicAssignment(fooTopicId, 0)))
)));
GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder()
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_ASSIGNORS_CONFIG, List.of(assignor))
.withMetadataImage(metadataImage)
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_MAX_SIZE_CONFIG, groupMaxSize)
.withConsumerGroup(new ConsumerGroupBuilder(groupId, 10)
.withMember(new ConsumerGroupMember.Builder(oldMemberId)
.setState(MemberState.STABLE)
.setInstanceId(instanceId)
.setMemberEpoch(10)
.setPreviousMemberEpoch(9)
.setRebalanceTimeoutMs(5000)
.setSubscribedTopicNames(List.of("foo", "bar"))
.setServerAssignorName("range")
.setAssignedPartitions(mkAssignment(
mkTopicAssignment(fooTopicId, 0)))
.build())
.withAssignment(oldMemberId, mkAssignment(mkTopicAssignment(fooTopicId, 0)))
.withAssignmentEpoch(10)
.withMetadataHash(computeGroupHash(Map.of(fooTopicName, computeTopicHash(fooTopicName, metadataImage)))))
.build();

assertThrows(UnreleasedInstanceIdException.class, () -> context.consumerGroupHeartbeat(
new ConsumerGroupHeartbeatRequestData()
.setGroupId(groupId)
.setInstanceId(instanceId)
.setMemberId(newMemberId)
.setMemberEpoch(0)
.setServerAssignor("range")
.setRebalanceTimeoutMs(5000)
.setSubscribedTopicNames(List.of("foo"))
.setTopicPartitions(List.of())));
}

@Test
public void testNewStaticMemberIsRejectedWhenConsumerGroupIsFull() {
String groupId = "fooup";
String instanceId = "instance-id";
String otherInstanceId = "other-instance-id";
String oldMemberId = "old-member-id";
String newMemberId = "new-member-id";
int groupMaxSize = 1;

Uuid fooTopicId = Uuid.randomUuid();
String fooTopicName = "foo";

CoordinatorMetadataImage metadataImage = new MetadataImageBuilder()
.addTopic(fooTopicId, fooTopicName, 1)
.buildCoordinatorMetadataImage();

MockPartitionAssignor assignor = new MockPartitionAssignor("range");
assignor.prepareGroupAssignment(new GroupAssignment(Map.of(
oldMemberId, new MemberAssignmentImpl(mkAssignment(mkTopicAssignment(fooTopicId, 0)))
)));
GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder()
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_ASSIGNORS_CONFIG, List.of(assignor))
.withMetadataImage(metadataImage)
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_MAX_SIZE_CONFIG, groupMaxSize)
.withConsumerGroup(new ConsumerGroupBuilder(groupId, 10)
.withMember(new ConsumerGroupMember.Builder(oldMemberId)
.setState(MemberState.STABLE)
.setInstanceId(instanceId)
.setMemberEpoch(10)
.setPreviousMemberEpoch(9)
.setRebalanceTimeoutMs(5000)
.setSubscribedTopicNames(List.of("foo", "bar"))
.setServerAssignorName("range")
.setAssignedPartitions(mkAssignment(
mkTopicAssignment(fooTopicId, 0)))
.build())
.withAssignment(oldMemberId, mkAssignment(mkTopicAssignment(fooTopicId, 0)))
.withAssignmentEpoch(10)
.withMetadataHash(computeGroupHash(Map.of(fooTopicName, computeTopicHash(fooTopicName, metadataImage)))))
.build();

assertThrows(GroupMaxSizeReachedException.class, () -> context.consumerGroupHeartbeat(
new ConsumerGroupHeartbeatRequestData()
.setGroupId(groupId)
.setInstanceId(otherInstanceId)
.setMemberId(newMemberId)
.setMemberEpoch(0)
.setServerAssignor("range")
.setRebalanceTimeoutMs(5000)
.setSubscribedTopicNames(List.of("foo"))
.setTopicPartitions(List.of())));
}

@Test
public void testConsumerGroupStates() {
String groupId = "fooup";
Expand Down Expand Up @@ -13724,6 +13895,74 @@ 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.UPGRADE.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)
.build();

JoinGroupRequestData request = new GroupMetadataManagerTestContext.JoinGroupRequestBuilder()
.withGroupId(groupId)
.withMemberId(UNKNOWN_MEMBER_ID)
.withGroupInstanceId(newInstanceId)
.withProtocols(GroupMetadataManagerTestContext.toConsumerProtocol(List.of(), List.of()))
.build();

assertThrows(GroupMaxSizeReachedException.class,
() -> context.sendClassicGroupJoin(request, true, true)
);
}


@Test
public void testJoiningConsumerGroupThrowsExceptionIfProtocolIsNotSupported() {
String groupId = "group-id";
Expand Down