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 31b9c4e6f94 KAFKA-20292 [4/N]: Lift member updates out of streams
TargetAssignmentBuilder (#22692)
31b9c4e6f94 is described below
commit 31b9c4e6f9457c7ca77535cc5c75d9bbd9430985
Author: Sean Quah <[email protected]>
AuthorDate: Mon Jun 29 12:56:25 2026 +0100
KAFKA-20292 [4/N]: Lift member updates out of streams
TargetAssignmentBuilder (#22692)
Remove the member update logic from the streams group
TargetAssignmentBuilder. Instead, use the new
UpdatedMembersAndTargetAssignmentView class to provide updated views of
members and the target assignment to the TargetAssignmentBuilder.
Reviewers: David Jacot <[email protected]>, Lucas Brutschy
<[email protected]>
---
.../coordinator/group/GroupMetadataManager.java | 27 +-
.../group/streams/TargetAssignmentBuilder.java | 74 ----
...sGroupStaticMemberGroupMetadataManagerTest.java | 1 -
.../group/streams/TargetAssignmentBuilderTest.java | 424 +--------------------
4 files changed, 13 insertions(+), 513 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 78b42edf754..c5d86de30c2 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
@@ -4333,6 +4333,16 @@ public class GroupMetadataManager {
TaskAssignor assignor = streamsGroupAssignor(group.groupId());
try {
+ UpdatedMembersAndTargetAssignmentView<StreamsGroupMember,
TasksTuple> updatedMembersAndTargetAssignment =
+ new UpdatedMembersAndTargetAssignmentView<>(
+ group.members(),
+ group.staticMembers(),
+ group.targetAssignment()
+ );
+ updatedMember.ifPresent(member ->
+
updatedMembersAndTargetAssignment.addOrUpdateMember(member.memberId(),
member.instanceId().orElse(null), member)
+ );
+
org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder
assignmentResultBuilder =
new
org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder(
group.groupId(),
@@ -4341,23 +4351,10 @@ public class GroupMetadataManager {
assignmentConfigs
)
.withTime(time)
- .withMembers(group.members())
+ .withMembers(updatedMembersAndTargetAssignment.members())
.withTopology(configuredTopology)
- .withStaticMembers(group.staticMembers())
.withMetadataImage(metadataImage)
- .withTargetAssignment(group.targetAssignment());
-
- updatedMember.ifPresent(member -> {
- assignmentResultBuilder.addOrUpdateMember(member.memberId(),
member);
- // If the instance id was associated to a different member, it
means that the
- // static member is replaced by the current member hence we
remove the previous one.
- member.instanceId().ifPresent(instanceId -> {
- StreamsGroupMember previousMember =
group.staticMember(instanceId);
- if (previousMember != null &&
!member.memberId().equals(previousMember.memberId())) {
-
assignmentResultBuilder.removeMember(previousMember.memberId());
- }
- });
- });
+
.withTargetAssignment(updatedMembersAndTargetAssignment.targetAssignment());
long startTimeMs = time.milliseconds();
org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder.TargetAssignmentResult
assignmentResult =
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TargetAssignmentBuilder.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TargetAssignmentBuilder.java
index f00f629dd08..072dbf90e6f 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TargetAssignmentBuilder.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TargetAssignmentBuilder.java
@@ -71,11 +71,6 @@ public class TargetAssignmentBuilder {
*/
private final Map<String, String> assignmentConfigs;
- /**
- * The members which have been updated or deleted. A null value signals
deleted members.
- */
- private final Map<String, StreamsGroupMember> updatedMembers = new
HashMap<>();
-
/**
* The members in the group.
*/
@@ -96,11 +91,6 @@ public class TargetAssignmentBuilder {
*/
private ConfiguredTopology topology;
- /**
- * The static members in the group.
- */
- private Map<String, String> staticMembers = Map.of();
-
/**
* Constructs the object.
*
@@ -161,19 +151,6 @@ public class TargetAssignmentBuilder {
return this;
}
- /**
- * Adds all the existing static members.
- *
- * @param staticMembers The existing static members in the streams group.
- * @return This object.
- */
- public TargetAssignmentBuilder withStaticMembers(
- Map<String, String> staticMembers
- ) {
- this.staticMembers = staticMembers;
- return this;
- }
-
/**
* Adds the metadata image to use.
*
@@ -213,33 +190,6 @@ public class TargetAssignmentBuilder {
return this;
}
- /**
- * Adds or updates a member. This is useful when the updated member is not
yet materialized in memory.
- *
- * @param memberId The member ID.
- * @param member The member to add or update.
- * @return This object.
- */
- public TargetAssignmentBuilder addOrUpdateMember(
- String memberId,
- StreamsGroupMember member
- ) {
- this.updatedMembers.put(memberId, member);
- return this;
- }
-
- /**
- * Removes a member. This is useful when the removed member is not yet
materialized in memory.
- *
- * @param memberId The member ID.
- * @return This object.
- */
- public TargetAssignmentBuilder removeMember(
- String memberId
- ) {
- return addOrUpdateMember(memberId, null);
- }
-
/**
* Builds the new target assignment.
*
@@ -255,30 +205,6 @@ public class TargetAssignmentBuilder {
targetAssignment.getOrDefault(memberId,
org.apache.kafka.coordinator.group.streams.TasksTuple.EMPTY)
)));
- // Update the member spec if updated or deleted members.
- updatedMembers.forEach((memberId, updatedMemberOrNull) -> {
- if (updatedMemberOrNull == null) {
- memberSpecs.remove(memberId);
- } else {
- org.apache.kafka.coordinator.group.streams.TasksTuple
assignment = targetAssignment.getOrDefault(memberId,
-
org.apache.kafka.coordinator.group.streams.TasksTuple.EMPTY);
-
- // A new static member joins and needs to replace an existing
departed one.
- if (updatedMemberOrNull.instanceId().isPresent()) {
- String previousMemberId =
staticMembers.get(updatedMemberOrNull.instanceId().get());
- if (previousMemberId != null &&
!previousMemberId.equals(memberId)) {
- assignment =
targetAssignment.getOrDefault(previousMemberId,
-
org.apache.kafka.coordinator.group.streams.TasksTuple.EMPTY);
- }
- }
-
- memberSpecs.put(memberId, createAssignmentMemberSpec(
- updatedMemberOrNull,
- assignment
- ));
- }
- });
-
// Compute the assignment.
GroupAssignment newGroupAssignment;
if (topology.isReady()) {
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 073a613b5d5..4890f19bb53 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
@@ -1444,7 +1444,6 @@ class StreamsGroupStaticMemberGroupMetadataManagerTest {
StreamsCoordinatorRecordHelpers.newStreamsGroupMemberRecord(groupId,
newJoinStaticMember),
StreamsCoordinatorRecordHelpers.newStreamsGroupMetadataRecord(groupId,
bumpedGroupEpoch, topic.metadataHash(), 0, getDefaultAssignmentConfigs(), -1,
-1),
-
StreamsCoordinatorRecordHelpers.newStreamsGroupTargetAssignmentRecord(groupId,
newJoinStaticMember.memberId(), givenTargetAssignment),
StreamsCoordinatorRecordHelpers.newStreamsGroupTargetAssignmentMetadataRecord(groupId,
bumpedGroupEpoch, context.time.milliseconds()),
StreamsCoordinatorRecordHelpers.newStreamsGroupCurrentAssignmentRecord(groupId,
reconciledMember)
),
diff --git
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/TargetAssignmentBuilderTest.java
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/TargetAssignmentBuilderTest.java
index 9e59a60d088..d4e00e51969 100644
---
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/TargetAssignmentBuilderTest.java
+++
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/TargetAssignmentBuilderTest.java
@@ -261,178 +261,6 @@ public class TargetAssignmentBuilderTest {
}
- @ParameterizedTest
- @EnumSource(TaskRole.class)
- public void testNewMember(TaskRole taskRole) {
- TargetAssignmentBuilderTestContext context = new
TargetAssignmentBuilderTestContext(
- "my-group",
- 20,
- 12345L
- );
-
- String fooSubtopologyId =
context.addSubtopologyWithSingleSourceTopic("foo", 6);
- String barSubtopologyId =
context.addSubtopologyWithSingleSourceTopic("bar", 6);
-
- context.addGroupMember("member-1", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2, 3),
- mkTasks(barSubtopologyId, 1, 2, 3)
- ));
-
- context.addGroupMember("member-2", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 4, 5, 6),
- mkTasks(barSubtopologyId, 4, 5, 6)
- ));
-
- context.updateMemberMetadata("member-3");
-
- context.prepareMemberAssignment("member-1", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2),
- mkTasks(barSubtopologyId, 1, 2)
- ));
-
- context.prepareMemberAssignment("member-2", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 3, 4),
- mkTasks(barSubtopologyId, 3, 4)
- ));
-
- context.prepareMemberAssignment("member-3", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 5, 6),
- mkTasks(barSubtopologyId, 5, 6)
- ));
-
-
org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder.TargetAssignmentResult
result = context.build();
-
- assertEquals(4, result.records().size());
-
- assertUnorderedRecordsEquals(List.of(List.of(
- newStreamsGroupTargetAssignmentRecord("my-group", "member-1",
mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2),
- mkTasks(barSubtopologyId, 1, 2)
- )),
- newStreamsGroupTargetAssignmentRecord("my-group", "member-2",
mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 3, 4),
- mkTasks(barSubtopologyId, 3, 4)
- )),
- newStreamsGroupTargetAssignmentRecord("my-group", "member-3",
mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 5, 6),
- mkTasks(barSubtopologyId, 5, 6)
- ))
- )), result.records().subList(0, 3));
-
- assertEquals(newStreamsGroupTargetAssignmentMetadataRecord(
- "my-group",
- 20,
- 12345L
- ), result.records().get(3));
-
- Map<String, TasksTuple> expectedAssignment = new HashMap<>();
- expectedAssignment.put("member-1", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2),
- mkTasks(barSubtopologyId, 1, 2)
- ));
- expectedAssignment.put("member-2", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 3, 4),
- mkTasks(barSubtopologyId, 3, 4)
- ));
- expectedAssignment.put("member-3", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 5, 6),
- mkTasks(barSubtopologyId, 5, 6)
- ));
-
- assertEquals(expectedAssignment, result.targetAssignment());
- }
-
-
- @ParameterizedTest
- @EnumSource(TaskRole.class)
- public void testUpdateMember(TaskRole taskRole) {
- TargetAssignmentBuilderTestContext context = new
TargetAssignmentBuilderTestContext(
- "my-group",
- 20,
- 12345L
- );
-
- String fooSubtopologyId =
context.addSubtopologyWithSingleSourceTopic("foo", 6);
- String barSubtopologyId =
context.addSubtopologyWithSingleSourceTopic("bar", 6);
-
- context.addGroupMember("member-1", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2, 3),
- mkTasks(barSubtopologyId, 1, 2)
- ));
-
- context.addGroupMember("member-2", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 4, 5, 6),
- mkTasks(barSubtopologyId, 3, 4)
- ));
-
- context.addGroupMember("member-3", mkTasksTuple(taskRole,
- mkTasks(barSubtopologyId, 5, 6)
- ));
-
- context.updateMemberMetadata(
- "member-3",
- Optional.of("instance-id-3"),
- Optional.of("rack-0")
- );
-
- context.prepareMemberAssignment("member-1", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2),
- mkTasks(barSubtopologyId, 1, 2)
- ));
-
- context.prepareMemberAssignment("member-2", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 3, 4),
- mkTasks(barSubtopologyId, 3, 4)
- ));
-
- context.prepareMemberAssignment("member-3", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 5, 6),
- mkTasks(barSubtopologyId, 5, 6)
- ));
-
-
org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder.TargetAssignmentResult
result = context.build();
-
- assertEquals(4, result.records().size());
-
- assertUnorderedRecordsEquals(List.of(List.of(
- newStreamsGroupTargetAssignmentRecord("my-group", "member-1",
mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2),
- mkTasks(barSubtopologyId, 1, 2)
- )),
- newStreamsGroupTargetAssignmentRecord("my-group", "member-2",
mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 3, 4),
- mkTasks(barSubtopologyId, 3, 4)
- )),
- newStreamsGroupTargetAssignmentRecord("my-group", "member-3",
mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 5, 6),
- mkTasks(barSubtopologyId, 5, 6)
- ))
- )), result.records().subList(0, 3));
-
- assertEquals(newStreamsGroupTargetAssignmentMetadataRecord(
- "my-group",
- 20,
- 12345L
- ), result.records().get(3));
-
- Map<String, TasksTuple> expectedAssignment = new HashMap<>();
- expectedAssignment.put("member-1", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2),
- mkTasks(barSubtopologyId, 1, 2)
- ));
- expectedAssignment.put("member-2", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 3, 4),
- mkTasks(barSubtopologyId, 3, 4)
- ));
- expectedAssignment.put("member-3", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 5, 6),
- mkTasks(barSubtopologyId, 5, 6)
- ));
-
- assertEquals(expectedAssignment, result.targetAssignment());
- }
-
-
@ParameterizedTest
@EnumSource(TaskRole.class)
public void testPartialAssignmentUpdate(TaskRole taskRole) {
@@ -515,164 +343,6 @@ public class TargetAssignmentBuilderTest {
}
- @ParameterizedTest
- @EnumSource(TaskRole.class)
- public void testDeleteMember(TaskRole taskRole) {
- TargetAssignmentBuilderTestContext context = new
TargetAssignmentBuilderTestContext(
- "my-group",
- 20,
- 12345L
- );
-
- String fooSubtopologyId =
context.addSubtopologyWithSingleSourceTopic("foo", 6);
- String barSubtopologyId =
context.addSubtopologyWithSingleSourceTopic("bar", 6);
-
- context.addGroupMember("member-1", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2),
- mkTasks(barSubtopologyId, 1, 2)
- ));
-
- context.addGroupMember("member-2", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 3, 4),
- mkTasks(barSubtopologyId, 3, 4)
- ));
-
- context.addGroupMember("member-3", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 5, 6),
- mkTasks(barSubtopologyId, 5, 6)
- ));
-
- context.removeMember("member-3");
-
- context.prepareMemberAssignment("member-1", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2, 3),
- mkTasks(barSubtopologyId, 1, 2, 3)
- ));
-
- context.prepareMemberAssignment("member-2", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 4, 5, 6),
- mkTasks(barSubtopologyId, 4, 5, 6)
- ));
-
-
org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder.TargetAssignmentResult
result = context.build();
-
- assertEquals(3, result.records().size());
-
- assertUnorderedRecordsEquals(List.of(List.of(
- newStreamsGroupTargetAssignmentRecord("my-group", "member-1",
mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2, 3),
- mkTasks(barSubtopologyId, 1, 2, 3)
- )),
- newStreamsGroupTargetAssignmentRecord("my-group", "member-2",
mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 4, 5, 6),
- mkTasks(barSubtopologyId, 4, 5, 6)
- ))
- )), result.records().subList(0, 2));
-
- assertEquals(newStreamsGroupTargetAssignmentMetadataRecord(
- "my-group",
- 20,
- 12345L
- ), result.records().get(2));
-
- Map<String, TasksTuple> expectedAssignment = new HashMap<>();
- expectedAssignment.put("member-1", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2, 3),
- mkTasks(barSubtopologyId, 1, 2, 3)
- ));
- expectedAssignment.put("member-2", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 4, 5, 6),
- mkTasks(barSubtopologyId, 4, 5, 6)
- ));
-
- assertEquals(expectedAssignment, result.targetAssignment());
- }
-
-
- @ParameterizedTest
- @EnumSource(TaskRole.class)
- public void testReplaceStaticMember(TaskRole taskRole) {
- TargetAssignmentBuilderTestContext context = new
TargetAssignmentBuilderTestContext(
- "my-group",
- 20,
- 12345L
- );
-
- String fooSubtopologyId =
context.addSubtopologyWithSingleSourceTopic("foo", 6);
- String barSubtopologyId =
context.addSubtopologyWithSingleSourceTopic("bar", 6);
-
- context.addGroupMember("member-1", "instance-member-1",
mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2),
- mkTasks(barSubtopologyId, 1, 2)
- ));
-
- context.addGroupMember("member-2", "instance-member-2",
mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 3, 4),
- mkTasks(barSubtopologyId, 3, 4)
- ));
-
- context.addGroupMember("member-3", "instance-member-3",
mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 5, 6),
- mkTasks(barSubtopologyId, 5, 6)
- ));
-
- // Static member 3 leaves
- context.removeMember("member-3");
-
- // Another static member joins with the same instance id as the
departed one
- context.updateMemberMetadata("member-3-a",
Optional.of("instance-member-3"),
- Optional.empty());
-
- context.prepareMemberAssignment("member-1", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2),
- mkTasks(barSubtopologyId, 1, 2)
- ));
-
- context.prepareMemberAssignment("member-2", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 3, 4),
- mkTasks(barSubtopologyId, 3, 4)
- ));
-
- context.prepareMemberAssignment("member-3-a", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 5, 6),
- mkTasks(barSubtopologyId, 5, 6)
- ));
-
- TargetAssignmentBuilder.TargetAssignmentResult result =
context.build();
-
- assertEquals(2, result.records().size());
-
- assertUnorderedRecordsEquals(List.of(List.of(
- newStreamsGroupTargetAssignmentRecord("my-group", "member-3-a",
mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 5, 6),
- mkTasks(barSubtopologyId, 5, 6)
- ))
- )), result.records().subList(0, 1));
-
- assertEquals(newStreamsGroupTargetAssignmentMetadataRecord(
- "my-group",
- 20,
- 12345L
- ), result.records().get(1));
-
- Map<String, TasksTuple> expectedAssignment = new HashMap<>();
- expectedAssignment.put("member-1", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 1, 2),
- mkTasks(barSubtopologyId, 1, 2)
- ));
- expectedAssignment.put("member-2", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 3, 4),
- mkTasks(barSubtopologyId, 3, 4)
- ));
-
- expectedAssignment.put("member-3-a", mkTasksTuple(taskRole,
- mkTasks(fooSubtopologyId, 5, 6),
- mkTasks(barSubtopologyId, 5, 6)
- ));
-
- assertEquals(expectedAssignment, result.targetAssignment());
- }
-
public static class TargetAssignmentBuilderTestContext {
private final String groupId;
@@ -684,10 +354,8 @@ public class TargetAssignmentBuilderTest {
Optional.empty());
private final Map<String, StreamsGroupMember> members = new
HashMap<>();
private final Map<String,
org.apache.kafka.coordinator.group.streams.TopicMetadata> subscriptionMetadata
= new HashMap<>();
- private final Map<String, StreamsGroupMember> updatedMembers = new
HashMap<>();
private final Map<String, TasksTuple> targetAssignment = new
HashMap<>();
private final Map<String, MemberAssignment> memberAssignments = new
HashMap<>();
- private final Map<String, String> staticMembers = new HashMap<>();
private MetadataImageBuilder topicsImageBuilder = new
MetadataImageBuilder();
public TargetAssignmentBuilderTestContext(
@@ -703,26 +371,12 @@ public class TargetAssignmentBuilderTest {
public void addGroupMember(
String memberId,
TasksTuple targetTasks
- ) {
- addGroupMember(memberId, null, targetTasks);
- }
-
- private void addGroupMember(
- String memberId,
- String instanceId,
- TasksTuple targetTasks
) {
StreamsGroupMember.Builder memberBuilder = new
StreamsGroupMember.Builder(memberId);
memberBuilder.setProcessId("processId");
memberBuilder.setClientTags(Map.of());
memberBuilder.setUserEndpoint(new
StreamsGroupMemberMetadataValue.Endpoint().setHost("host").setPort(9090));
-
- if (instanceId != null) {
- memberBuilder.setInstanceId(instanceId);
- staticMembers.put(instanceId, memberId);
- } else {
- memberBuilder.setInstanceId(null);
- }
+ memberBuilder.setInstanceId(null);
memberBuilder.setRackId(null);
members.put(memberId, memberBuilder.build());
targetAssignment.put(memberId, targetTasks);
@@ -740,45 +394,6 @@ public class TargetAssignmentBuilderTest {
return subtopologyId;
}
- public void updateMemberMetadata(
- String memberId
- ) {
- updateMemberMetadata(
- memberId,
- Optional.empty(),
- Optional.empty()
- );
- }
-
- public void updateMemberMetadata(
- String memberId,
- Optional<String> instanceId,
- Optional<String> rackId
- ) {
- StreamsGroupMember existingMember = members.get(memberId);
- StreamsGroupMember.Builder builder;
- if (existingMember != null) {
- builder = new StreamsGroupMember.Builder(existingMember);
- } else {
- builder = new StreamsGroupMember.Builder(memberId);
- builder.setProcessId("processId");
- builder.setRackId(null);
- builder.setInstanceId(null);
- builder.setClientTags(Map.of());
- builder.setUserEndpoint(new
StreamsGroupMemberMetadataValue.Endpoint().setHost("host").setPort(9090));
- }
- updatedMembers.put(memberId, builder
- .maybeUpdateInstanceId(instanceId)
- .maybeUpdateRackId(rackId)
- .build());
- }
-
- public void removeMember(
- String memberId
- ) {
- this.updatedMembers.put(memberId, null);
- }
-
public void prepareMemberAssignment(
String memberId,
TasksTuple assignment
@@ -789,8 +404,6 @@ public class TargetAssignmentBuilderTest {
public
org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder.TargetAssignmentResult
build() {
// Prepare expected member specs.
Map<String, AssignmentMemberSpec> memberSpecs = new HashMap<>();
-
- // All the existing members are prepared.
members.forEach((memberId, member) ->
memberSpecs.put(memberId, createAssignmentMemberSpec(
member,
@@ -798,31 +411,6 @@ public class TargetAssignmentBuilderTest {
)
));
- // All the updated are added and all the deleted
- // members are removed.
- updatedMembers.forEach((memberId, updatedMemberOrNull) -> {
- if (updatedMemberOrNull == null) {
- memberSpecs.remove(memberId);
- } else {
- TasksTuple assignment =
targetAssignment.getOrDefault(memberId,
- TasksTuple.EMPTY);
-
- // A new static member joins and needs to replace an
existing departed one.
- if (updatedMemberOrNull.instanceId().isPresent()) {
- String previousMemberId =
staticMembers.get(updatedMemberOrNull.instanceId().get());
- if (previousMemberId != null &&
!previousMemberId.equals(memberId)) {
- assignment =
targetAssignment.getOrDefault(previousMemberId,
- TasksTuple.EMPTY);
- }
- }
-
- memberSpecs.put(memberId, createAssignmentMemberSpec(
- updatedMemberOrNull,
- assignment
- ));
- }
- });
-
CoordinatorMetadataImage metadataImage = new
KRaftCoordinatorMetadataImage(topicsImageBuilder.build());
// Prepare the expected topology metadata.
@@ -842,19 +430,9 @@ public class TargetAssignmentBuilderTest {
.withTime(new MockTime(0, assignmentTimestamp,
assignmentTimestamp))
.withMembers(members)
.withTopology(topology)
- .withStaticMembers(staticMembers)
.withMetadataImage(metadataImage)
.withTargetAssignment(targetAssignment);
- // Add the updated members or delete the deleted members.
- updatedMembers.forEach((memberId, updatedMemberOrNull) -> {
- if (updatedMemberOrNull != null) {
- builder.addOrUpdateMember(memberId, updatedMemberOrNull);
- } else {
- builder.removeMember(memberId);
- }
- });
-
// Execute the builder.
org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder.TargetAssignmentResult
result = builder.build();