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)