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 4626eccfae7 KAFKA-20292 [5/N]: Fix 
UpdatedMembersAndTargetAssignmentView to handle instance id changes (#22713)
4626eccfae7 is described below

commit 4626eccfae79c1459f38bfd624f4fb1f2b0966ab
Author: Sean Quah <[email protected]>
AuthorDate: Wed Jul 1 20:17:18 2026 +0100

    KAFKA-20292 [5/N]: Fix UpdatedMembersAndTargetAssignmentView to handle 
instance id changes (#22713)
    
    Fix addOrUpdateMember in UpdatedMembersAndTargetAssignmentView to remove
    a member's previous instance id mapping when its instance id changes.
    Previously, the method 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]>
---
 .../coordinator/group/GroupMetadataManager.java    |  15 +--
 .../UpdatedMembersAndTargetAssignmentView.java     |  38 ++++++--
 .../UpdatedMembersAndTargetAssignmentViewTest.java | 107 +++++++++++++--------
 .../assignor/TargetAssignmentBuilderBenchmark.java |   4 +-
 4 files changed, 108 insertions(+), 56 deletions(-)

diff --git 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
index ad669498d5b..238e637e05c 100644
--- 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
+++ 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
@@ -4161,9 +4161,10 @@ public class GroupMetadataManager {
                 new UpdatedMembersAndTargetAssignmentView<>(
                     group.members(),
                     group.staticMembers(),
-                    group.targetAssignment()
+                    group.targetAssignment(),
+                    ConsumerGroupMember::instanceId
                 );
-            
updatedMembersAndTargetAssignment.addOrUpdateMember(updatedMember.memberId(), 
updatedMember.instanceId(), updatedMember);
+            
updatedMembersAndTargetAssignment.addOrUpdateMember(updatedMember.memberId(), 
updatedMember);
 
             TargetAssignmentBuilder.ConsumerTargetAssignmentBuilder 
assignmentResultBuilder =
                 new 
TargetAssignmentBuilder.ConsumerTargetAssignmentBuilder(group.groupId(), 
groupEpoch, consumerGroupAssignors.get(preferredServerAssignor))
@@ -4244,9 +4245,10 @@ public class GroupMetadataManager {
                 new UpdatedMembersAndTargetAssignmentView<>(
                     group.members(),
                     Map.of(),
-                    group.targetAssignment()
+                    group.targetAssignment(),
+                    ShareGroupMember::instanceId
                 );
-            
updatedMembersAndTargetAssignment.addOrUpdateMember(updatedMember.memberId(), 
updatedMember.instanceId(), updatedMember);
+            
updatedMembersAndTargetAssignment.addOrUpdateMember(updatedMember.memberId(), 
updatedMember);
 
             TargetAssignmentBuilder.ShareTargetAssignmentBuilder 
assignmentResultBuilder =
                 new 
TargetAssignmentBuilder.ShareTargetAssignmentBuilder(group.groupId(), 
groupEpoch, shareGroupAssignor)
@@ -4348,10 +4350,11 @@ public class GroupMetadataManager {
                 new UpdatedMembersAndTargetAssignmentView<>(
                     group.members(),
                     group.staticMembers(),
-                    group.targetAssignment()
+                    group.targetAssignment(),
+                    m -> m.instanceId().orElse(null)
                 );
             updatedMember.ifPresent(member ->
-                
updatedMembersAndTargetAssignment.addOrUpdateMember(member.memberId(), 
member.instanceId().orElse(null), member)
+                
updatedMembersAndTargetAssignment.addOrUpdateMember(member.memberId(), member)
             );
 
             org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder 
assignmentResultBuilder =
diff --git 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/util/UpdatedMembersAndTargetAssignmentView.java
 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/util/UpdatedMembersAndTargetAssignmentView.java
index 5494c69854e..962809007b8 100644
--- 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/util/UpdatedMembersAndTargetAssignmentView.java
+++ 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/util/UpdatedMembersAndTargetAssignmentView.java
@@ -19,6 +19,7 @@ package org.apache.kafka.coordinator.group.util;
 import java.util.Collections;
 import java.util.Map;
 import java.util.Objects;
+import java.util.function.Function;
 
 /**
  * A view of a group's members, static members, and target assignment after 
unwritten membership
@@ -44,19 +45,27 @@ public class UpdatedMembersAndTargetAssignmentView<M, A> {
      */
     private final OverlayMap<String, A> targetAssignment;
 
+    /**
+     * Gets a member's instance id, or {@code null} if the member is not 
static.
+     */
+    private final Function<M, String> getInstanceId;
+
     /**
      * @param members          The group members. Must not be modified during 
the lifetime of the view.
      * @param staticMembers    The static group members. Must not be modified 
during the lifetime of the view.
      * @param targetAssignment The target assignment per member id. Must not 
be modified during the lifetime of the view.
+     * @param getInstanceId    Gets a member's instance id, or {@code null} if 
the member is not static.
      */
     public UpdatedMembersAndTargetAssignmentView(
         Map<String, M> members,
         Map<String, String> staticMembers,
-        Map<String, A> targetAssignment
+        Map<String, A> targetAssignment,
+        Function<M, String> getInstanceId
     ) {
         this.members = new OverlayMap<>(Objects.requireNonNull(members));
         this.staticMembers = new 
OverlayMap<>(Objects.requireNonNull(staticMembers));
         this.targetAssignment = new 
OverlayMap<>(Objects.requireNonNull(targetAssignment));
+        this.getInstanceId = Objects.requireNonNull(getInstanceId);
     }
 
     /**
@@ -85,12 +94,21 @@ public class UpdatedMembersAndTargetAssignmentView<M, A> {
      * member for the same instance id, the previous static member's target 
assignment is moved to
      * the new member and the previous static member is removed from the view.
      *
-     * @param memberId   The member id.
-     * @param instanceId The instance id of the member, or {@code null} if the 
member is not static.
-     * @param member     The member to add or update.
+     * @param memberId The member id.
+     * @param member   The member to add or update.
      */
-    public void addOrUpdateMember(String memberId, String instanceId, M 
member) {
-        members.put(memberId, member);
+    public void addOrUpdateMember(String memberId, M member) {
+        M previousMember = members.put(memberId, member);
+        String previousInstanceId = previousMember != null ? 
getInstanceId.apply(previousMember) : null;
+        String instanceId = getInstanceId.apply(member);
+
+        // Remove the old static member mapping when the instance id has 
changed.
+        // We don't remove the mapping when the instance id has not changed, 
otherwise we won't
+        // detect static member replacement correctly below.
+        if (previousInstanceId != null && 
!previousInstanceId.equals(instanceId)) {
+            staticMembers.remove(previousInstanceId);
+        }
+
         if (instanceId != null) {
             String previousMemberId = staticMembers.put(instanceId, memberId);
             if (previousMemberId != null && 
!memberId.equals(previousMemberId)) {
@@ -110,11 +128,11 @@ public class UpdatedMembersAndTargetAssignmentView<M, A> {
     /**
      * Removes a member.
      *
-     * @param memberId   The member id.
-     * @param instanceId The instance id of the member, or {@code null} if the 
member is not static.
+     * @param memberId The member id.
      */
-    public void removeMember(String memberId, String instanceId) {
-        members.remove(memberId);
+    public void removeMember(String memberId) {
+        M member = members.remove(memberId);
+        String instanceId = member != null ? getInstanceId.apply(member) : 
null;
         if (instanceId != null && 
memberId.equals(staticMembers.get(instanceId))) {
             staticMembers.remove(instanceId);
         }
diff --git 
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/util/UpdatedMembersAndTargetAssignmentViewTest.java
 
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/util/UpdatedMembersAndTargetAssignmentViewTest.java
index 84be50efca0..843fc3f5c8a 100644
--- 
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/util/UpdatedMembersAndTargetAssignmentViewTest.java
+++ 
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/util/UpdatedMembersAndTargetAssignmentViewTest.java
@@ -24,15 +24,20 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
 
 public class UpdatedMembersAndTargetAssignmentViewTest {
 
+    /**
+     * A test member.
+     */
+    private record Member(String name, String instanceId) { }
+
     /**
      * Creates an {@link UpdatedMembersAndTargetAssignmentView} with two 
members, one static and one
      * non-static.
      */
-    private static UpdatedMembersAndTargetAssignmentView<String, String> 
createView() {
+    private static UpdatedMembersAndTargetAssignmentView<Member, String> 
createView() {
         return new UpdatedMembersAndTargetAssignmentView<>(
             Map.of(
-                "member-1", "Member1",
-                "member-2", "Member2"
+                "member-1", new Member("Member1", null),
+                "member-2", new Member("Member2", "instance-id")
             ),
             Map.of(
                 "instance-id", "member-2"
@@ -40,20 +45,21 @@ public class UpdatedMembersAndTargetAssignmentViewTest {
             Map.of(
                 "member-1", "Assignment-member-1",
                 "member-2", "Assignment-member-2"
-            )
+            ),
+            Member::instanceId
         );
     }
 
     @Test
     public void testAddMember() {
-        UpdatedMembersAndTargetAssignmentView<String, String> view = 
createView();
+        UpdatedMembersAndTargetAssignmentView<Member, String> view = 
createView();
 
-        view.addOrUpdateMember("member-3", null, "Member3");
+        view.addOrUpdateMember("member-3", new Member("Member3", null));
 
         assertEquals(Map.of(
-            "member-1", "Member1",
-            "member-2", "Member2",
-            "member-3", "Member3"
+            "member-1", new Member("Member1", null),
+            "member-2", new Member("Member2", "instance-id"),
+            "member-3", new Member("Member3", null)
         ), view.members());
         assertEquals(Map.of(
             "instance-id", "member-2"
@@ -66,14 +72,14 @@ public class UpdatedMembersAndTargetAssignmentViewTest {
 
     @Test
     public void testAddStaticMember() {
-        UpdatedMembersAndTargetAssignmentView<String, String> view = 
createView();
+        UpdatedMembersAndTargetAssignmentView<Member, String> view = 
createView();
 
-        view.addOrUpdateMember("member-3", "instance-id-2", "Member3");
+        view.addOrUpdateMember("member-3", new Member("Member3", 
"instance-id-2"));
 
         assertEquals(Map.of(
-            "member-1", "Member1",
-            "member-2", "Member2",
-            "member-3", "Member3"
+            "member-1", new Member("Member1", null),
+            "member-2", new Member("Member2", "instance-id"),
+            "member-3", new Member("Member3", "instance-id-2")
         ), view.members());
         assertEquals(Map.of(
             "instance-id", "member-2",
@@ -87,13 +93,13 @@ public class UpdatedMembersAndTargetAssignmentViewTest {
 
     @Test
     public void testReplaceMember() {
-        UpdatedMembersAndTargetAssignmentView<String, String> view = 
createView();
+        UpdatedMembersAndTargetAssignmentView<Member, String> view = 
createView();
 
-        view.addOrUpdateMember("member-1", null, "Member1-updated");
+        view.addOrUpdateMember("member-1", new Member("Member1-updated", 
null));
 
         assertEquals(Map.of(
-            "member-1", "Member1-updated",
-            "member-2", "Member2"
+            "member-1", new Member("Member1-updated", null),
+            "member-2", new Member("Member2", "instance-id")
         ), view.members());
         assertEquals(Map.of(
             "instance-id", "member-2"
@@ -106,13 +112,13 @@ public class UpdatedMembersAndTargetAssignmentViewTest {
 
     @Test
     public void testReplaceStaticMemberWithSameMemberId() {
-        UpdatedMembersAndTargetAssignmentView<String, String> view = 
createView();
+        UpdatedMembersAndTargetAssignmentView<Member, String> view = 
createView();
 
-        view.addOrUpdateMember("member-2", "instance-id", "Member2-updated");
+        view.addOrUpdateMember("member-2", new Member("Member2-updated", 
"instance-id"));
 
         assertEquals(Map.of(
-            "member-1", "Member1",
-            "member-2", "Member2-updated"
+            "member-1", new Member("Member1", null),
+            "member-2", new Member("Member2-updated", "instance-id")
         ), view.members());
         assertEquals(Map.of(
             "instance-id", "member-2"
@@ -125,13 +131,13 @@ public class UpdatedMembersAndTargetAssignmentViewTest {
 
     @Test
     public void testReplaceStaticMemberWithDifferentMemberId() {
-        UpdatedMembersAndTargetAssignmentView<String, String> view = 
createView();
+        UpdatedMembersAndTargetAssignmentView<Member, String> view = 
createView();
 
-        view.addOrUpdateMember("member-3", "instance-id", "Member3");
+        view.addOrUpdateMember("member-3", new Member("Member3", 
"instance-id"));
 
         assertEquals(Map.of(
-            "member-1", "Member1",
-            "member-3", "Member3"
+            "member-1", new Member("Member1", null),
+            "member-3", new Member("Member3", "instance-id")
         ), view.members());
         assertEquals(Map.of(
             "instance-id", "member-3"
@@ -140,31 +146,56 @@ public class UpdatedMembersAndTargetAssignmentViewTest {
             "member-1", "Assignment-member-1",
             "member-3", "Assignment-member-2"
         ), view.targetAssignment());
+    }
 
-        // Removing the previous static member does not change the new static 
member's assignment.
-        view.removeMember("member-2", "instance-id");
+    @Test
+    public void testReplaceStaticMemberWithNullInstanceId() {
+        // This operation is not possible, since a heartbeat with a null 
instance id will keep any
+        // existing instance id.
+        UpdatedMembersAndTargetAssignmentView<Member, String> view = 
createView();
+
+        view.addOrUpdateMember("member-2", new Member("Member2-updated", 
null));
 
         assertEquals(Map.of(
-            "member-1", "Member1",
-            "member-3", "Member3"
+            "member-1", new Member("Member1", null),
+            "member-2", new Member("Member2-updated", null)
         ), view.members());
+        assertEquals(Map.of(), view.staticMembers());
         assertEquals(Map.of(
-            "instance-id", "member-3"
+            "member-1", "Assignment-member-1",
+            "member-2", "Assignment-member-2"
+        ), view.targetAssignment());
+    }
+
+    @Test
+    public void testReplaceStaticMemberWithDifferentInstanceId() {
+        // This operation will never happen with the official Java client and 
may be forbidden in
+        // the future.
+        UpdatedMembersAndTargetAssignmentView<Member, String> view = 
createView();
+
+        view.addOrUpdateMember("member-2", new Member("Member2-updated", 
"instance-id-2"));
+
+        assertEquals(Map.of(
+            "member-1", new Member("Member1", null),
+            "member-2", new Member("Member2-updated", "instance-id-2")
+        ), view.members());
+        assertEquals(Map.of(
+            "instance-id-2", "member-2"
         ), view.staticMembers());
         assertEquals(Map.of(
             "member-1", "Assignment-member-1",
-            "member-3", "Assignment-member-2"
+            "member-2", "Assignment-member-2"
         ), view.targetAssignment());
     }
 
     @Test
     public void testRemoveMember() {
-        UpdatedMembersAndTargetAssignmentView<String, String> view = 
createView();
+        UpdatedMembersAndTargetAssignmentView<Member, String> view = 
createView();
 
-        view.removeMember("member-1", null);
+        view.removeMember("member-1");
 
         assertEquals(Map.of(
-            "member-2", "Member2"
+            "member-2", new Member("Member2", "instance-id")
         ), view.members());
         assertEquals(Map.of(
             "instance-id", "member-2"
@@ -176,12 +207,12 @@ public class UpdatedMembersAndTargetAssignmentViewTest {
 
     @Test
     public void testRemoveStaticMember() {
-        UpdatedMembersAndTargetAssignmentView<String, String> view = 
createView();
+        UpdatedMembersAndTargetAssignmentView<Member, String> view = 
createView();
 
-        view.removeMember("member-2", "instance-id");
+        view.removeMember("member-2");
 
         assertEquals(Map.of(
-            "member-1", "Member1"
+            "member-1", new Member("Member1", null)
         ), view.members());
         assertEquals(Map.of(), view.staticMembers());
         assertEquals(Map.of(
diff --git 
a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/TargetAssignmentBuilderBenchmark.java
 
b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/TargetAssignmentBuilderBenchmark.java
index 2859b2f3961..35cc73e50c3 100644
--- 
a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/TargetAssignmentBuilderBenchmark.java
+++ 
b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/TargetAssignmentBuilderBenchmark.java
@@ -116,8 +116,8 @@ public class TargetAssignmentBuilderBenchmark {
             .setSubscribedTopicNames(allTopicNames)
             .build();
 
-        updatedMembersAndTargetAssignment = new 
UpdatedMembersAndTargetAssignmentView<>(members, Map.of(), 
existingTargetAssignment);
-        
updatedMembersAndTargetAssignment.addOrUpdateMember(newMember.memberId(), 
newMember.instanceId(), newMember);
+        updatedMembersAndTargetAssignment = new 
UpdatedMembersAndTargetAssignmentView<>(members, Map.of(), 
existingTargetAssignment, ConsumerGroupMember::instanceId);
+        
updatedMembersAndTargetAssignment.addOrUpdateMember(newMember.memberId(), 
newMember);
 
         targetAssignmentBuilder = new 
TargetAssignmentBuilder.ConsumerTargetAssignmentBuilder(GROUP_ID, GROUP_EPOCH, 
partitionAssignor)
             .withTime(Time.SYSTEM)

Reply via email to