This is an automated email from the ASF dual-hosted git repository.
squah-confluent pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 4c6a34e7b7c MINOR: Add static member replacement test with new
subscription and not assignment (#22698)
4c6a34e7b7c is described below
commit 4c6a34e7b7cfb10f7aad9dc6046d031c2f6dd7c0
Author: Sean Quah <[email protected]>
AuthorDate: Tue Jun 30 08:41:29 2026 +0100
MINOR: Add static member replacement test with new subscription and not
assignment (#22698)
There is an existing test for a consumer group static member rejoining
with an updated subscription and an updated assignment. Add a version of
the test where the new assignment is the same as the old one. Prior to
#22510, we would emit a redundant target assignment record for the
static member.
Reviewers: David Jacot <[email protected]>
---
.../group/GroupMetadataManagerTest.java | 190 +++++++++++++++++++++
1 file changed, 190 insertions(+)
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 48652f4a73b..d4218bc5adf 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
@@ -2541,6 +2541,196 @@ public class GroupMetadataManagerTest {
context.assertNoRebalanceTimeout(groupId, memberId2);
}
+ @Test
+ public void
testStaticMemberRejoinsWithNewSubscribedTopicsAndNoAssignmentChange() {
+ String groupId = "fooup";
+ // Use a static member id as it makes the test easier.
+ String memberId1 = Uuid.randomUuid().toString();
+ String memberId2 = Uuid.randomUuid().toString();
+ String member2RejoinId = Uuid.randomUuid().toString();
+
+ Uuid fooTopicId = Uuid.randomUuid();
+ String fooTopicName = "foo";
+ Uuid barTopicId = Uuid.randomUuid();
+ String barTopicName = "bar";
+
+ MockPartitionAssignor assignor = new MockPartitionAssignor("range");
+ ConsumerGroupMember member1 = new
ConsumerGroupMember.Builder(memberId1)
+ .setState(MemberState.STABLE)
+ .setInstanceId("instance-id-1")
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(9)
+ .setRebalanceTimeoutMs(5000)
+ .setClientId(DEFAULT_CLIENT_ID)
+ .setClientHost(DEFAULT_CLIENT_ADDRESS.toString())
+ .setSubscribedTopicNames(List.of("foo", "bar"))
+ .setServerAssignorName("range")
+ .setAssignedPartitions(toAssignmentWithEpochs(mkAssignment(
+ mkTopicAssignment(fooTopicId, 0, 1, 2),
+ mkTopicAssignment(barTopicId, 0)), 10))
+ .build();
+ ConsumerGroupMember member2 = new
ConsumerGroupMember.Builder(memberId2)
+ .setState(MemberState.STABLE)
+ .setInstanceId("instance-id-2")
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(9)
+ .setRebalanceTimeoutMs(5000)
+ .setClientId(DEFAULT_CLIENT_ID)
+ .setClientHost(DEFAULT_CLIENT_ADDRESS.toString())
+ .setSubscribedTopicNames(List.of("foo"))
+ .setServerAssignorName("range")
+ .setAssignedPartitions(toAssignmentWithEpochs(mkAssignment(
+ mkTopicAssignment(fooTopicId, 3, 4, 5)), 10))
+ .build();
+
+ MetadataImage metadataImage = new MetadataImageBuilder()
+ .addTopic(fooTopicId, fooTopicName, 6)
+ .addTopic(barTopicId, barTopicName, 3)
+ .addRacks()
+ .build();
+ long fooTopicHash = computeTopicHash(fooTopicName, new
KRaftCoordinatorMetadataImage(metadataImage));
+ long barTopicHash = computeTopicHash(barTopicName, new
KRaftCoordinatorMetadataImage(metadataImage));
+
+ // Consumer group with two static members.
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_ASSIGNORS_CONFIG,
List.of(assignor))
+ .withMetadataImage(new
KRaftCoordinatorMetadataImage(metadataImage))
+ .withConsumerGroup(new ConsumerGroupBuilder(groupId, 10)
+ .withMember(member1)
+ .withMember(member2)
+ .withAssignment(memberId1, mkAssignment(
+ mkTopicAssignment(fooTopicId, 0, 1, 2),
+ mkTopicAssignment(barTopicId, 0)))
+ .withAssignment(memberId2, mkAssignment(
+ mkTopicAssignment(fooTopicId, 3, 4, 5)))
+ .withAssignmentEpoch(10)
+ .withMetadataHash(computeGroupHash(Map.of(
+ fooTopicName, fooTopicHash,
+ barTopicName, barTopicHash
+ ))))
+ .build();
+
+ assignor.prepareGroupAssignment(new GroupAssignment(Map.of(
+ memberId1, new MemberAssignmentImpl(mkAssignment(
+ mkTopicAssignment(fooTopicId, 0, 1, 2),
+ mkTopicAssignment(barTopicId, 0)
+ )),
+ member2RejoinId, new MemberAssignmentImpl(mkAssignment(
+ mkTopicAssignment(fooTopicId, 3, 4, 5)
+ ))
+ )));
+
+ // Member 2 leaves the consumer group.
+ CoordinatorResult<ConsumerGroupHeartbeatResponseData,
CoordinatorRecord> result = context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setGroupId(groupId)
+ .setMemberId(memberId2)
+ .setInstanceId("instance-id-2")
+ .setMemberEpoch(-2));
+
+ // Member epoch of the response would be set to -2.
+ assertResponseEquals(
+ new ConsumerGroupHeartbeatResponseData()
+ .setMemberId(memberId2)
+ .setMemberEpoch(-2),
+ result.response()
+ );
+
+ // The departing static member will have it's epoch set to -2.
+ ConsumerGroupMember member2UpdatedEpoch = new
ConsumerGroupMember.Builder(member2)
+ .setMemberEpoch(-2)
+ .setPartitionsPendingRevocation(Map.of())
+ .resetAssignedPartitionsEpochsToZero()
+ .build();
+
+ assertEquals(1, result.records().size());
+ assertRecordEquals(result.records().get(0),
GroupCoordinatorRecordHelpers.newConsumerGroupCurrentAssignmentRecord(groupId,
member2UpdatedEpoch));
+
+ // Member 2 rejoins the group with the same instance id.
+ CoordinatorResult<ConsumerGroupHeartbeatResponseData,
CoordinatorRecord> rejoinResult = context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setMemberId(member2RejoinId)
+ .setGroupId(groupId)
+ .setInstanceId("instance-id-2")
+ .setMemberEpoch(0)
+ .setRebalanceTimeoutMs(5000)
+ .setServerAssignor("range")
+ .setSubscribedTopicNames(List.of("foo", "bar")) // bar is new.
+ .setTopicPartitions(List.of()));
+
+ assertResponseEquals(
+ new ConsumerGroupHeartbeatResponseData()
+ .setMemberId(member2RejoinId)
+ .setMemberEpoch(11)
+ .setHeartbeatIntervalMs(5000)
+ .setAssignment(new
ConsumerGroupHeartbeatResponseData.Assignment()
+ .setTopicPartitions(List.of(
+ new
ConsumerGroupHeartbeatResponseData.TopicPartitions()
+ .setTopicId(fooTopicId)
+ .setPartitions(List.of(3, 4, 5))
+ ))),
+ rejoinResult.response()
+ );
+
+ ConsumerGroupMember expectedCopiedMember = new
ConsumerGroupMember.Builder(member2RejoinId)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(0)
+ .setPreviousMemberEpoch(0)
+ .setInstanceId("instance-id-2")
+ .setClientId(DEFAULT_CLIENT_ID)
+ .setClientHost(DEFAULT_CLIENT_ADDRESS.toString())
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(List.of("foo"))
+ .setServerAssignorName("range")
+ .setAssignedPartitions(toAssignmentWithEpochs(mkAssignment(
+ mkTopicAssignment(fooTopicId, 3, 4, 5)), 0))
+ .build();
+
+ ConsumerGroupMember expectedRejoinedMember = new
ConsumerGroupMember.Builder(member2RejoinId)
+ .setState(MemberState.STABLE)
+ .setMemberEpoch(11)
+ .setPreviousMemberEpoch(0)
+ .setInstanceId("instance-id-2")
+ .setClientId(DEFAULT_CLIENT_ID)
+ .setClientHost(DEFAULT_CLIENT_ADDRESS.toString())
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(List.of("foo", "bar"))
+ .setServerAssignorName("range")
+ .setAssignedPartitions(mkAssignmentWithEpochs(
+ mkTopicAssignmentWithEpochs(fooTopicId, 0, 3, 4, 5)
+ ))
+ .build();
+
+ List<CoordinatorRecord> expectedRecordsAfterRejoin = List.of(
+ // The previous member is deleted.
+
GroupCoordinatorRecordHelpers.newConsumerGroupCurrentAssignmentTombstoneRecord(groupId,
memberId2),
+
GroupCoordinatorRecordHelpers.newConsumerGroupTargetAssignmentTombstoneRecord(groupId,
memberId2),
+
GroupCoordinatorRecordHelpers.newConsumerGroupMemberSubscriptionTombstoneRecord(groupId,
memberId2),
+
+ // The new member is created as a copy of the previous one but
+ // with its new member id and new epochs.
+
GroupCoordinatorRecordHelpers.newConsumerGroupMemberSubscriptionRecord(groupId,
expectedCopiedMember),
+
GroupCoordinatorRecordHelpers.newConsumerGroupTargetAssignmentRecord(groupId,
member2RejoinId, mkAssignment(
+ mkTopicAssignment(fooTopicId, 3, 4, 5))),
+
GroupCoordinatorRecordHelpers.newConsumerGroupCurrentAssignmentRecord(groupId,
expectedCopiedMember),
+
+ // As the new member as a different subscribed topic set, a
rebalance is triggered.
+
GroupCoordinatorRecordHelpers.newConsumerGroupMemberSubscriptionRecord(groupId,
expectedRejoinedMember),
+ GroupCoordinatorRecordHelpers.newConsumerGroupEpochRecord(groupId,
11, computeGroupHash(Map.of(
+ fooTopicName, fooTopicHash,
+ barTopicName, barTopicHash
+ ))),
+ // The target assignment is unchanged, so no new target assignment
record is written.
+
GroupCoordinatorRecordHelpers.newConsumerGroupTargetAssignmentMetadataRecord(groupId,
11, context.time.milliseconds()),
+
GroupCoordinatorRecordHelpers.newConsumerGroupCurrentAssignmentRecord(groupId,
expectedRejoinedMember)
+ );
+
+ assertRecordsEquals(expectedRecordsAfterRejoin,
rejoinResult.records());
+ // Verify that there are no timers.
+ context.assertNoSessionTimeout(groupId, memberId2);
+ context.assertNoRebalanceTimeout(groupId, memberId2);
+ }
+
@Test
public void testStaticMembersRejoinWithNewServerAssignor() {
String groupId = "fooup";