diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml
index a317fa558c10f..eaeb29d3243a2 100644
--- a/checkstyle/suppressions.xml
+++ b/checkstyle/suppressions.xml
@@ -347,7 +347,7 @@
+ files="(GroupMetadataManager|GroupMetadataManagerTest).java"/>
records, String groupId, Strin
private void cancelTimers(String groupId, String memberId) {
cancelConsumerGroupSessionTimeout(groupId, memberId);
cancelConsumerGroupRebalanceTimeout(groupId, memberId);
+ cancelConsumerGroupJoinTimeout(groupId, memberId);
cancelConsumerGroupSyncTimeout(groupId, memberId);
}
@@ -3985,10 +3986,17 @@ public CoordinatorResult classicGroupSync(
RequestContext context,
SyncGroupRequestData request,
CompletableFuture responseFuture
- ) throws UnknownMemberIdException, GroupIdNotFoundException {
- Group group = groups.get(request.groupId(), Long.MAX_VALUE);
+ ) throws UnknownMemberIdException {
+ Group group;
+ try {
+ group = group(request.groupId());
+ } catch (GroupIdNotFoundException e) {
+ responseFuture.complete(new SyncGroupResponseData()
+ .setErrorCode(Errors.UNKNOWN_MEMBER_ID.code()));
+ return EMPTY_RESULT;
+ }
- if (group == null || group.isEmpty()) {
+ if (group.isEmpty()) {
responseFuture.complete(new SyncGroupResponseData()
.setErrorCode(Errors.UNKNOWN_MEMBER_ID.code()));
return EMPTY_RESULT;
@@ -4250,9 +4258,10 @@ public CoordinatorResult classicGroupH
RequestContext context,
HeartbeatRequestData request
) {
- Group group = groups.get(request.groupId(), Long.MAX_VALUE);
-
- if (group == null) {
+ Group group;
+ try {
+ group = group(request.groupId());
+ } catch (GroupIdNotFoundException e) {
throw new UnknownMemberIdException(
String.format("Group %s not found.", request.groupId())
);
@@ -4426,14 +4435,129 @@ private ConsumerGroupMember validateConsumerGroupMember(
* @param context The request context.
* @param request The actual LeaveGroup request.
*
+ * @return The LeaveGroup response and the records to append.
+ */
+ public CoordinatorResult classicGroupLeave(
+ RequestContext context,
+ LeaveGroupRequestData request
+ ) throws UnknownMemberIdException {
+ Group group;
+ try {
+ group = group(request.groupId());
+ } catch (GroupIdNotFoundException e) {
+ throw new UnknownMemberIdException(String.format("Group %s not found.", request.groupId()));
+ }
+
+ if (group.type() == CLASSIC) {
+ return classicGroupLeaveToClassicGroup((ClassicGroup) group, context, request);
+ } else {
+ return classicGroupLeaveToConsumerGroup((ConsumerGroup) group, context, request);
+ }
+ }
+
+ /**
+ * Handle a classic LeaveGroupRequest to a ConsumerGroup.
+ *
+ * @param group The ConsumerGroup.
+ * @param context The request context.
+ * @param request The actual LeaveGroup request.
+ *
+ * @return The LeaveGroup response and the records to append.
+ */
+ private CoordinatorResult classicGroupLeaveToConsumerGroup(
+ ConsumerGroup group,
+ RequestContext context,
+ LeaveGroupRequestData request
+ ) throws UnknownMemberIdException {
+ String groupId = group.groupId();
+ List memberResponses = new ArrayList<>();
+ Set validLeaveGroupMembers = new HashSet<>();
+ List records = new ArrayList<>();
+
+ for (MemberIdentity memberIdentity : request.members()) {
+ String memberId = memberIdentity.memberId();
+ String instanceId = memberIdentity.groupInstanceId();
+ String reason = memberIdentity.reason() != null ? memberIdentity.reason() : "not provided";
+
+ ConsumerGroupMember member;
+ try {
+ if (instanceId == null) {
+ member = group.getOrMaybeCreateMember(memberId, false);
+ throwIfMemberDoesNotUseClassicProtocol(member);
+
+ log.info("[Group {}] Dynamic Member {} has left group " +
+ "through explicit `LeaveGroup` request; client reason: {}",
+ groupId, memberId, reason);
+ } else {
+ member = group.staticMember(instanceId);
+ throwIfStaticMemberIsUnknown(member, instanceId);
+ // The LeaveGroup API allows administrative removal of members by GroupInstanceId
+ // in which case we expect the MemberId to be undefined.
+ if (!UNKNOWN_MEMBER_ID.equals(memberId)) {
+ throwIfInstanceIdIsFenced(member, groupId, memberId, instanceId);
+ }
+ throwIfMemberDoesNotUseClassicProtocol(member);
+
+ memberId = member.memberId();
+ log.info("[Group {}] Static Member {} with instance id {} has left group " +
+ "through explicit `LeaveGroup` request; client reason: {}",
+ groupId, memberId, instanceId, reason);
+ }
+
+ removeMember(records, groupId, memberId);
+ cancelTimers(groupId, memberId);
+ memberResponses.add(
+ new MemberResponse()
+ .setMemberId(memberId)
+ .setGroupInstanceId(instanceId)
+ );
+ validLeaveGroupMembers.add(member);
+ } catch (KafkaException e) {
+ memberResponses.add(
+ new MemberResponse()
+ .setMemberId(memberId)
+ .setGroupInstanceId(instanceId)
+ .setErrorCode(Errors.forException(e).code())
+ );
+ }
+ }
+
+ if (!records.isEmpty()) {
+ // Maybe update the subscription metadata.
+ Map subscriptionMetadata = group.computeSubscriptionMetadata(
+ group.computeSubscribedTopicNames(validLeaveGroupMembers),
+ metadataImage.topics(),
+ metadataImage.cluster()
+ );
+
+ if (!subscriptionMetadata.equals(group.subscriptionMetadata())) {
+ log.info("[GroupId {}] Computed new subscription metadata: {}.",
+ group.groupId(), subscriptionMetadata);
+ records.add(newGroupSubscriptionMetadataRecord(group.groupId(), subscriptionMetadata));
+ }
+
+ // Bump the group epoch.
+ records.add(newGroupEpochRecord(groupId, group.groupEpoch() + 1));
+ }
+
+ return new CoordinatorResult<>(records, new LeaveGroupResponseData().setMembers(memberResponses));
+ }
+
+ /**
+ * Handle a classic LeaveGroupRequest to a ClassicGroup.
+ *
+ * @param group The ClassicGroup.
+ * @param context The request context.
+ * @param request The actual LeaveGroup request.
+ *
* @return The LeaveGroup response and the GroupMetadata record to append if the group
* no longer has any members.
*/
- public CoordinatorResult classicGroupLeave(
+ private CoordinatorResult classicGroupLeaveToClassicGroup(
+ ClassicGroup group,
RequestContext context,
LeaveGroupRequestData request
- ) throws UnknownMemberIdException, GroupIdNotFoundException {
- ClassicGroup group = getOrMaybeCreateClassicGroup(request.groupId(), false);
+ ) throws UnknownMemberIdException {
if (group.isInState(DEAD)) {
return new CoordinatorResult<>(
Collections.emptyList(),
diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroup.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroup.java
index 007971a95dd24..bbc544289b26b 100644
--- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroup.java
+++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroup.java
@@ -1076,6 +1076,29 @@ public Map computeSubscribedTopicNames(
return subscribedTopicNames;
}
+ /**
+ * Updates the subscription count with a set of members removed.
+ *
+ * @param removedMembers The set of removed members.
+ *
+ * @return Copy of the map of topics to the count of number of subscribers.
+ */
+ public Map computeSubscribedTopicNames(
+ Set removedMembers
+ ) {
+ Map subscribedTopicNames = new HashMap<>(this.subscribedTopicNames);
+ if (removedMembers != null) {
+ removedMembers.forEach(removedMember ->
+ maybeUpdateSubscribedTopicNames(
+ subscribedTopicNames,
+ removedMember,
+ null
+ )
+ );
+ }
+ return subscribedTopicNames;
+ }
+
/**
* Compute the subscription type of the consumer group.
*
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 3fa2cbb0aee9d..aa90e07c5dcf2 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
@@ -12801,6 +12801,365 @@ public void testConsumerGroupMemberUsingClassicProtocolFencedWhenJoinTimeout() {
);
}
+ @Test
+ public void testConsumerGroupMemberUsingClassicProtocolBatchLeaveGroup() {
+ String groupId = "group-id";
+ String memberId1 = Uuid.randomUuid().toString();
+ String memberId2 = Uuid.randomUuid().toString();
+ String memberId3 = Uuid.randomUuid().toString();
+ String instanceId2 = "instance-id-2";
+ String instanceId3 = "instance-id-3";
+
+ Uuid fooTopicId = Uuid.randomUuid();
+ String fooTopicName = "foo";
+ Uuid barTopicId = Uuid.randomUuid();
+ String barTopicName = "bar";
+
+ List protocol1 = Collections.singletonList(
+ new ConsumerGroupMemberMetadataValue.ClassicProtocol()
+ .setName("range")
+ .setMetadata(Utils.toArray(ConsumerProtocol.serializeSubscription(new ConsumerPartitionAssignor.Subscription(
+ Arrays.asList(fooTopicName, barTopicName),
+ null,
+ Collections.singletonList(new TopicPartition(fooTopicName, 0))
+ ))))
+ );
+ List protocol2 = Collections.singletonList(
+ new ConsumerGroupMemberMetadataValue.ClassicProtocol()
+ .setName("range")
+ .setMetadata(Utils.toArray(ConsumerProtocol.serializeSubscription(new ConsumerPartitionAssignor.Subscription(
+ Arrays.asList(fooTopicName, barTopicName),
+ null,
+ Collections.singletonList(new TopicPartition(fooTopicName, 1))
+ ))))
+ );
+
+ ConsumerGroupMember member1 = new ConsumerGroupMember.Builder(memberId1)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(9)
+ .setClientId("client")
+ .setClientHost("localhost/127.0.0.1")
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setServerAssignorName("range")
+ .setRebalanceTimeoutMs(45000)
+ .setClassicMemberMetadata(
+ new ConsumerGroupMemberMetadataValue.ClassicMemberMetadata()
+ .setSessionTimeoutMs(5000)
+ .setSupportedProtocols(protocol1)
+ )
+ .setAssignedPartitions(mkAssignment(mkTopicAssignment(fooTopicId, 0)))
+ .build();
+ ConsumerGroupMember member2 = new ConsumerGroupMember.Builder(memberId2)
+ .setInstanceId(instanceId2)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(9)
+ .setPreviousMemberEpoch(8)
+ .setClientId("client")
+ .setClientHost("localhost/127.0.0.1")
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setServerAssignorName("range")
+ .setRebalanceTimeoutMs(45000)
+ .setClassicMemberMetadata(
+ new ConsumerGroupMemberMetadataValue.ClassicMemberMetadata()
+ .setSessionTimeoutMs(5000)
+ .setSupportedProtocols(protocol2)
+ )
+ .setAssignedPartitions(mkAssignment(mkTopicAssignment(fooTopicId, 1)))
+ .build();
+ ConsumerGroupMember member3 = new ConsumerGroupMember.Builder(memberId3)
+ .setInstanceId(instanceId3)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(9)
+ .setClientId("client")
+ .setClientHost("localhost/127.0.0.1")
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setServerAssignorName("range")
+ .setRebalanceTimeoutMs(45000)
+ .setAssignedPartitions(mkAssignment(mkTopicAssignment(barTopicId, 0)))
+ .build();
+
+ // Consumer group with three members.
+ // Dynamic member 1 uses the classic protocol.
+ // Static member 2 uses the classic protocol.
+ // Static member 3 uses the consumer protocol.
+ GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(new MockPartitionAssignor("range")))
+ .withMetadataImage(new MetadataImageBuilder()
+ .addTopic(fooTopicId, fooTopicName, 2)
+ .addTopic(barTopicId, barTopicName, 1)
+ .addRacks()
+ .build())
+ .withConsumerGroup(new ConsumerGroupBuilder(groupId, 10)
+ .withMember(member1)
+ .withMember(member2)
+ .withMember(member3)
+ .withAssignment(memberId1, mkAssignment(mkTopicAssignment(fooTopicId, 0)))
+ .withAssignment(memberId2, mkAssignment(mkTopicAssignment(fooTopicId, 1)))
+ .withAssignment(memberId3, mkAssignment(mkTopicAssignment(barTopicId, 0)))
+ .withAssignmentEpoch(10))
+ .build();
+ context.groupMetadataManager.consumerGroup(groupId).setMetadataRefreshDeadline(Long.MAX_VALUE, 10);
+ context.replay(CoordinatorRecordHelpers.newGroupSubscriptionMetadataRecord(groupId, new HashMap() {
+ {
+ put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 2, mkMapOfPartitionRacks(2)));
+ put(barTopicName, new TopicMetadata(barTopicId, barTopicName, 1, mkMapOfPartitionRacks(1)));
+ }
+ }));
+
+ // Member 1 joins to schedule the sync timeout and the heartbeat timeout.
+ context.sendClassicGroupJoin(
+ new GroupMetadataManagerTestContext.JoinGroupRequestBuilder()
+ .withGroupId(groupId)
+ .withMemberId(memberId1)
+ .withRebalanceTimeoutMs(member1.rebalanceTimeoutMs())
+ .withSessionTimeoutMs(member1.classicMemberMetadata().get().sessionTimeoutMs())
+ .withProtocols(GroupMetadataManagerTestContext.toConsumerProtocol(
+ Arrays.asList(fooTopicName, barTopicName),
+ Collections.singletonList(new TopicPartition(fooTopicName, 0))))
+ .build()
+ ).appendFuture.complete(null);
+ context.assertSyncTimeout(groupId, memberId1, member1.rebalanceTimeoutMs());
+ context.assertSessionTimeout(groupId, memberId1, member1.classicMemberMetadata().get().sessionTimeoutMs());
+
+ // Member 2 heartbeats to schedule the join timeout and the heartbeat timeout.
+ context.sendClassicGroupHeartbeat(
+ new HeartbeatRequestData()
+ .setGroupId(groupId)
+ .setMemberId(memberId2)
+ .setGenerationId(9)
+ );
+ context.assertJoinTimeout(groupId, memberId2, member2.rebalanceTimeoutMs());
+ context.assertSessionTimeout(groupId, memberId2, member2.classicMemberMetadata().get().sessionTimeoutMs());
+
+ // Member 1 and member 2 leave the group.
+ CoordinatorResult leaveResult = context.sendClassicGroupLeave(
+ new LeaveGroupRequestData()
+ .setGroupId("group-id")
+ .setMembers(Arrays.asList(
+ // Valid member id.
+ new MemberIdentity()
+ .setMemberId(memberId1),
+ new MemberIdentity()
+ .setGroupInstanceId(instanceId2),
+ // Member that doesn't use the classic protocol.
+ new MemberIdentity()
+ .setMemberId(memberId3)
+ .setGroupInstanceId(instanceId3),
+ // Unknown member id.
+ new MemberIdentity()
+ .setMemberId("unknown-member-id"),
+ new MemberIdentity()
+ .setGroupInstanceId("unknown-instance-id"),
+ // Fenced instance id.
+ new MemberIdentity()
+ .setMemberId("unknown-member-id")
+ .setGroupInstanceId(instanceId3)
+ ))
+ );
+
+ assertEquals(
+ new LeaveGroupResponseData()
+ .setMembers(Arrays.asList(
+ new LeaveGroupResponseData.MemberResponse()
+ .setGroupInstanceId(null)
+ .setMemberId(memberId1),
+ new LeaveGroupResponseData.MemberResponse()
+ .setGroupInstanceId(instanceId2)
+ .setMemberId(memberId2),
+ new LeaveGroupResponseData.MemberResponse()
+ .setGroupInstanceId(instanceId3)
+ .setMemberId(memberId3)
+ .setErrorCode(Errors.UNKNOWN_MEMBER_ID.code()),
+ new LeaveGroupResponseData.MemberResponse()
+ .setGroupInstanceId(null)
+ .setMemberId("unknown-member-id")
+ .setErrorCode(Errors.UNKNOWN_MEMBER_ID.code()),
+ new LeaveGroupResponseData.MemberResponse()
+ .setGroupInstanceId("unknown-instance-id")
+ .setErrorCode(Errors.UNKNOWN_MEMBER_ID.code()),
+ new LeaveGroupResponseData.MemberResponse()
+ .setGroupInstanceId(instanceId3)
+ .setMemberId("unknown-member-id")
+ .setErrorCode(Errors.FENCED_INSTANCE_ID.code())
+ )),
+ leaveResult.response()
+ );
+
+ List expectedRecords = Arrays.asList(
+ // Remove member 1
+ CoordinatorRecordHelpers.newCurrentAssignmentTombstoneRecord(groupId, memberId1),
+ CoordinatorRecordHelpers.newTargetAssignmentTombstoneRecord(groupId, memberId1),
+ CoordinatorRecordHelpers.newMemberSubscriptionTombstoneRecord(groupId, memberId1),
+ // Remove member 2.
+ CoordinatorRecordHelpers.newCurrentAssignmentTombstoneRecord(groupId, memberId2),
+ CoordinatorRecordHelpers.newTargetAssignmentTombstoneRecord(groupId, memberId2),
+ CoordinatorRecordHelpers.newMemberSubscriptionTombstoneRecord(groupId, memberId2),
+ // Bump the group epoch.
+ CoordinatorRecordHelpers.newGroupEpochRecord(groupId, 11)
+ );
+ assertEquals(expectedRecords, leaveResult.records());
+
+ context.assertNoSessionTimeout(groupId, memberId1);
+ context.assertNoSyncTimeout(groupId, memberId1);
+ context.assertNoSessionTimeout(groupId, memberId2);
+ context.assertNoJoinTimeout(groupId, memberId2);
+ }
+
+ @Test
+ public void testConsumerGroupMemberUsingClassicProtocolBatchLeaveGroupUpdatingSubscriptionMetadata() {
+ String groupId = "group-id";
+ String memberId1 = Uuid.randomUuid().toString();
+ String memberId2 = Uuid.randomUuid().toString();
+
+ Uuid fooTopicId = Uuid.randomUuid();
+ String fooTopicName = "foo";
+ Uuid barTopicId = Uuid.randomUuid();
+ String barTopicName = "bar";
+
+ List protocol = Collections.singletonList(
+ new ConsumerGroupMemberMetadataValue.ClassicProtocol()
+ .setName("range")
+ .setMetadata(Utils.toArray(ConsumerProtocol.serializeSubscription(new ConsumerPartitionAssignor.Subscription(
+ Arrays.asList(fooTopicName, barTopicName),
+ null,
+ Collections.singletonList(new TopicPartition(fooTopicName, 0))
+ ))))
+ );
+
+ ConsumerGroupMember member1 = new ConsumerGroupMember.Builder(memberId1)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(9)
+ .setClientId("client")
+ .setClientHost("localhost/127.0.0.1")
+ .setSubscribedTopicNames(Arrays.asList("foo", "bar"))
+ .setServerAssignorName("range")
+ .setRebalanceTimeoutMs(45000)
+ .setClassicMemberMetadata(
+ new ConsumerGroupMemberMetadataValue.ClassicMemberMetadata()
+ .setSessionTimeoutMs(5000)
+ .setSupportedProtocols(protocol)
+ )
+ .setAssignedPartitions(mkAssignment(mkTopicAssignment(fooTopicId, 0)))
+ .build();
+ ConsumerGroupMember member2 = new ConsumerGroupMember.Builder(memberId2)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(9)
+ .setClientId("client")
+ .setClientHost("localhost/127.0.0.1")
+ .setSubscribedTopicNames(Arrays.asList("foo"))
+ .setServerAssignorName("range")
+ .setRebalanceTimeoutMs(45000)
+ .setAssignedPartitions(mkAssignment(mkTopicAssignment(barTopicId, 0)))
+ .build();
+
+ // Consumer group with two members.
+ // Member 1 uses the classic protocol and member 2 uses the consumer protocol.
+ GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(new MockPartitionAssignor("range")))
+ .withMetadataImage(new MetadataImageBuilder()
+ .addTopic(fooTopicId, fooTopicName, 2)
+ .addTopic(barTopicId, barTopicName, 1)
+ .addRacks()
+ .build())
+ .withConsumerGroup(new ConsumerGroupBuilder(groupId, 10)
+ .withMember(member1)
+ .withMember(member2)
+ .withAssignment(memberId1, mkAssignment(mkTopicAssignment(fooTopicId, 0)))
+ .withAssignment(memberId2, mkAssignment(mkTopicAssignment(barTopicId, 0)))
+ .withAssignmentEpoch(10))
+ .build();
+ context.groupMetadataManager.consumerGroup(groupId).setMetadataRefreshDeadline(Long.MAX_VALUE, 10);
+ context.replay(CoordinatorRecordHelpers.newGroupSubscriptionMetadataRecord(groupId, new HashMap() {
+ {
+ put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 2, mkMapOfPartitionRacks(2)));
+ put(barTopicName, new TopicMetadata(barTopicId, barTopicName, 1, mkMapOfPartitionRacks(1)));
+ }
+ }));
+
+ // Member 1 leaves the group.
+ CoordinatorResult leaveResult = context.sendClassicGroupLeave(
+ new LeaveGroupRequestData()
+ .setGroupId("group-id")
+ .setMembers(Collections.singletonList(
+ new MemberIdentity()
+ .setMemberId(memberId1)
+ ))
+ );
+
+ assertEquals(
+ new LeaveGroupResponseData()
+ .setMembers(Collections.singletonList(
+ new LeaveGroupResponseData.MemberResponse()
+ .setGroupInstanceId(null)
+ .setMemberId(memberId1)
+ )),
+ leaveResult.response()
+ );
+
+ List expectedRecords = Arrays.asList(
+ // Remove member 1
+ CoordinatorRecordHelpers.newCurrentAssignmentTombstoneRecord(groupId, memberId1),
+ CoordinatorRecordHelpers.newTargetAssignmentTombstoneRecord(groupId, memberId1),
+ CoordinatorRecordHelpers.newMemberSubscriptionTombstoneRecord(groupId, memberId1),
+ // Update the subscription metadata.
+ CoordinatorRecordHelpers.newGroupSubscriptionMetadataRecord(groupId, new HashMap() {
+ {
+ put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 2, mkMapOfPartitionRacks(2)));
+ }
+ }),
+ // Bump the group epoch.
+ CoordinatorRecordHelpers.newGroupEpochRecord(groupId, 11)
+ );
+ assertEquals(expectedRecords, leaveResult.records());
+ }
+
+ @Test
+ public void testClassicGroupLeaveToConsumerGroupWithoutValidLeaveGroupMember() {
+ String groupId = "group-id";
+ String memberId = Uuid.randomUuid().toString();
+
+ // Consumer group without member using the classic protocol.
+ GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder()
+ .withAssignors(Collections.singletonList(new MockPartitionAssignor("range")))
+ .withConsumerGroup(new ConsumerGroupBuilder(groupId, 10)
+ .withMember(new ConsumerGroupMember.Builder(memberId)
+ .build()))
+ .build();
+
+ // Send leave request without valid member.
+ CoordinatorResult leaveResult = context.sendClassicGroupLeave(
+ new LeaveGroupRequestData()
+ .setGroupId("group-id")
+ .setMembers(Arrays.asList(
+ new MemberIdentity()
+ .setMemberId("unknown-member-id"),
+ new MemberIdentity()
+ .setMemberId(memberId)
+ ))
+ );
+
+ assertEquals(
+ new LeaveGroupResponseData()
+ .setMembers(Arrays.asList(
+ new LeaveGroupResponseData.MemberResponse()
+ .setGroupInstanceId(null)
+ .setMemberId("unknown-member-id")
+ .setErrorCode(Errors.UNKNOWN_MEMBER_ID.code()),
+ new LeaveGroupResponseData.MemberResponse()
+ .setGroupInstanceId(null)
+ .setMemberId(memberId)
+ .setErrorCode(Errors.UNKNOWN_MEMBER_ID.code())
+ )),
+ leaveResult.response()
+ );
+
+ assertEquals(Collections.emptyList(), leaveResult.records());
+ }
+
private static void checkJoinGroupResponse(
JoinGroupResponseData expectedResponse,
JoinGroupResponseData actualResponse,
diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroupTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroupTest.java
index 784bb6077346b..a483fa3be022c 100644
--- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroupTest.java
+++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroupTest.java
@@ -765,6 +765,55 @@ public void testUpdateSubscriptionMetadata() {
image.cluster()
)
);
+
+ // Compute while taking into account removal of member 1, member 2 and member 3
+ assertEquals(
+ Collections.emptyMap(),
+ consumerGroup.computeSubscriptionMetadata(
+ consumerGroup.computeSubscribedTopicNames(new HashSet<>(Arrays.asList(member1, member2, member3))),
+ image.topics(),
+ image.cluster()
+ )
+ );
+
+ // Compute while taking into account removal of member 2 and member 3.
+ assertEquals(
+ mkMap(
+ mkEntry("foo", new TopicMetadata(fooTopicId, "foo", 1, mkMapOfPartitionRacks(1)))
+ ),
+ consumerGroup.computeSubscriptionMetadata(
+ consumerGroup.computeSubscribedTopicNames(new HashSet<>(Arrays.asList(member2, member3))),
+ image.topics(),
+ image.cluster()
+ )
+ );
+
+ // Compute while taking into account removal of member 1.
+ assertEquals(
+ mkMap(
+ mkEntry("bar", new TopicMetadata(barTopicId, "bar", 2, mkMapOfPartitionRacks(2))),
+ mkEntry("zar", new TopicMetadata(zarTopicId, "zar", 3, mkMapOfPartitionRacks(3)))
+ ),
+ consumerGroup.computeSubscriptionMetadata(
+ consumerGroup.computeSubscribedTopicNames(Collections.singleton(member1)),
+ image.topics(),
+ image.cluster()
+ )
+ );
+
+ // It should return foo, bar and zar.
+ assertEquals(
+ mkMap(
+ mkEntry("foo", new TopicMetadata(fooTopicId, "foo", 1, mkMapOfPartitionRacks(1))),
+ mkEntry("bar", new TopicMetadata(barTopicId, "bar", 2, mkMapOfPartitionRacks(2))),
+ mkEntry("zar", new TopicMetadata(zarTopicId, "zar", 3, mkMapOfPartitionRacks(3)))
+ ),
+ consumerGroup.computeSubscriptionMetadata(
+ consumerGroup.computeSubscribedTopicNames(Collections.emptySet()),
+ image.topics(),
+ image.cluster()
+ )
+ );
}
@Test