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 71faeec8b04 KAFKA-20747: Remove stale staticMembers entry when a
member changes its instance id (#22708)
71faeec8b04 is described below
commit 71faeec8b04ffd5f02b4154ca6f05a073641e303
Author: Sean Quah <[email protected]>
AuthorDate: Wed Jul 1 15:09:45 2026 +0100
KAFKA-20747: Remove stale staticMembers entry when a member changes its
instance id (#22708)
Fix updateStaticMember in ConsumerGroup and StreamsGroup to remove a
member's previous instance id mapping when its instance id changes.
Previously these methods only added the new instanceId -> memberId
mapping, so when a static member rejoined with the same member id but a
different instance id, the old mapping was left behind.
Reviewers: David Jacot <[email protected]>
---
.../group/modern/consumer/ConsumerGroup.java | 8 ++-
.../coordinator/group/streams/StreamsGroup.java | 9 +++-
.../group/GroupMetadataManagerTest.java | 63 ++++++++++++++++++++++
...sGroupStaticMemberGroupMetadataManagerTest.java | 48 +++++++++++++++++
4 files changed, 124 insertions(+), 4 deletions(-)
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/modern/consumer/ConsumerGroup.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/modern/consumer/ConsumerGroup.java
index 523f8e1873b..06f26e016d3 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/modern/consumer/ConsumerGroup.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/modern/consumer/ConsumerGroup.java
@@ -324,7 +324,7 @@ public class ConsumerGroup extends
ModernGroup<ConsumerGroupMember> {
maybeUpdateServerAssignors(oldMember, newMember);
maybeUpdatePartitionEpoch(oldMember, newMember);
maybeUpdateSubscribedRegularExpression(oldMember, newMember);
- updateStaticMember(newMember);
+ updateStaticMember(oldMember, newMember);
maybeUpdateGroupState();
maybeUpdateGroupSubscriptionType();
maybeUpdateNumClassicProtocolMembers(oldMember, newMember);
@@ -334,9 +334,13 @@ public class ConsumerGroup extends
ModernGroup<ConsumerGroupMember> {
/**
* Updates the member id stored against the instance id if the member is a
static member.
*
+ * @param oldMember The old member state.
* @param newMember The new member state.
*/
- private void updateStaticMember(ConsumerGroupMember newMember) {
+ private void updateStaticMember(ConsumerGroupMember oldMember,
ConsumerGroupMember newMember) {
+ if (oldMember != null && oldMember.instanceId() != null) {
+ staticMembers.remove(oldMember.instanceId());
+ }
if (newMember.instanceId() != null) {
staticMembers.put(newMember.instanceId(), newMember.memberId());
}
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/StreamsGroup.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/StreamsGroup.java
index 5f392e5a698..c9d57cb9e53 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/StreamsGroup.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/StreamsGroup.java
@@ -504,7 +504,7 @@ public class StreamsGroup implements Group {
}
StreamsGroupMember oldMember = members.put(newMember.memberId(),
newMember);
maybeUpdateTaskProcessId(oldMember, newMember);
- updateStaticMember(newMember);
+ updateStaticMember(oldMember, newMember);
maybeUpdateGroupState();
endpointToPartitionsCache.remove(newMember.memberId());
}
@@ -512,9 +512,14 @@ public class StreamsGroup implements Group {
/**
* Updates the member ID stored against the instance ID if the member is a
static member.
*
+ * @param oldMember The old member state.
* @param newMember The new member state.
*/
- private void updateStaticMember(StreamsGroupMember newMember) {
+ private void updateStaticMember(StreamsGroupMember oldMember,
StreamsGroupMember newMember) {
+ if (oldMember != null && oldMember.instanceId() != null &&
+ oldMember.instanceId().isPresent()) {
+ staticMembers.remove(oldMember.instanceId().get());
+ }
if (newMember.instanceId() != null &&
newMember.instanceId().isPresent()) {
staticMembers.put(newMember.instanceId().get(),
newMember.memberId());
}
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 d4218bc5adf..ef816ee2ac6 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
@@ -3104,6 +3104,69 @@ public class GroupMetadataManagerTest {
assertRecordsEquals(expectedRecords, result.records());
}
+ @Test
+ public void
testStaticMemberRejoinsWithSameMemberIdAndDifferentInstanceId() {
+ String groupId = "fooup";
+ String memberId1 = Uuid.randomUuid().toString();
+ String instanceId1 = "instance-1";
+ String instanceId2 = "instance-2";
+
+ Uuid fooTopicId = Uuid.randomUuid();
+ String fooTopicName = "foo";
+
+ MockPartitionAssignor assignor = new MockPartitionAssignor("range");
+
+ CoordinatorMetadataImage metadataImage = new MetadataImageBuilder()
+ .addTopic(fooTopicId, fooTopicName, 6)
+ .buildCoordinatorMetadataImage();
+
+ // Consumer group with one static member using instanceId1.
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+
.withConfig(GroupCoordinatorConfig.CONSUMER_GROUP_ASSIGNORS_CONFIG,
List.of(assignor))
+ .withMetadataImage(metadataImage)
+ .withConsumerGroup(new ConsumerGroupBuilder(groupId, 10)
+ .withMember(new ConsumerGroupMember.Builder(memberId1)
+ .setState(MemberState.STABLE)
+ .setInstanceId(instanceId1)
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(9)
+ .setClientId(DEFAULT_CLIENT_ID)
+ .setClientHost(DEFAULT_CLIENT_ADDRESS.toString())
+ .setSubscribedTopicNames(List.of("foo", "bar"))
+ .setServerAssignorName("range")
+ .setAssignedPartitions(toAssignmentWithEpochs(mkAssignment(
+ mkTopicAssignment(fooTopicId, 0, 1, 2, 3, 4, 5)), 10))
+ .build())
+ .withAssignment(memberId1, mkAssignment(
+ mkTopicAssignment(fooTopicId, 0, 1, 2, 3, 4, 5)))
+ .withAssignmentEpoch(10)
+ .withMetadataHash(computeGroupHash(Map.of(
+ fooTopicName, computeTopicHash(fooTopicName,
metadataImage))
+ )))
+ .build();
+
+ assertEquals(
+ Map.of(instanceId1, memberId1),
+ context.groupMetadataManager.consumerGroup(groupId).staticMembers()
+ );
+
+ // Member rejoins with the same member id and a different instance id.
+ context.consumerGroupHeartbeat(
+ new ConsumerGroupHeartbeatRequestData()
+ .setGroupId(groupId)
+ .setMemberId(memberId1)
+ .setInstanceId(instanceId2)
+ .setMemberEpoch(0)
+ .setRebalanceTimeoutMs(5000)
+ .setSubscribedTopicNames(List.of("foo", "bar"))
+ .setTopicPartitions(List.of()));
+
+ assertEquals(
+ Map.of(instanceId2, memberId1),
+ context.groupMetadataManager.consumerGroup(groupId).staticMembers()
+ );
+ }
+
@Test
public void
testShouldThrownUnreleasedInstanceIdExceptionWhenNewMemberJoinsWithInUseInstanceId()
{
String groupId = "fooup";
diff --git
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/StreamsGroupStaticMemberGroupMetadataManagerTest.java
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/StreamsGroupStaticMemberGroupMetadataManagerTest.java
index 4890f19bb53..3584313d7ca 100644
---
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/StreamsGroupStaticMemberGroupMetadataManagerTest.java
+++
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/StreamsGroupStaticMemberGroupMetadataManagerTest.java
@@ -1518,6 +1518,54 @@ class StreamsGroupStaticMemberGroupMetadataManagerTest {
));
}
+ @Test
+ public void
testStaticMemberRejoinsWithSameMemberIdAndDifferentInstanceId() {
+ int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+ String instanceId1 = "instance-1";
+ String instanceId2 = "instance-2";
+
+ StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo",
4);
+ TasksTuple targetAssignment = topic.targetAssignment(0, 1, 2, 3);
+ TasksTupleWithEpochs assignedTasks = topic.assignedTasks(groupEpoch,
0, 1, 2, 3);
+
+ // Streams group with one static member using instanceId1.
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(new
MockTaskAssignor("sticky")))
+ .withMetadataImage(topic.metadataImage())
+ .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId,
instanceId1)
+ .setMemberEpoch(groupEpoch)
+ .setPreviousMemberEpoch(groupEpoch - 1)
+ .setAssignedTasks(assignedTasks)
+ .build())
+ .withTargetAssignment(memberId, targetAssignment)
+ .withTargetAssignmentEpoch(groupEpoch)
+
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+ .withValidatedTopologyEpoch(0)
+ .withMetadataHash(topic.metadataHash())
+ .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+ .build();
+
+ assertEquals(
+ Map.of(instanceId1, memberId),
+ context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+ );
+
+ // Member rejoins with the same member id and a different instance id.
+ context.streamsGroupHeartbeat(
+ staticJoinHeartbeat(groupId, memberId, instanceId2, topic)
+ );
+
+ assertEquals(
+ Map.of(instanceId2, memberId),
+ context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+ );
+ }
+
@Test
public void
testStaticMemberLeaveWithMismatchedMemberIdThrowsFencedInstanceIdInStreamsGroup()
{
String groupId = "fooup";