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 6fa3711c3e9 KAFKA-20292 [3/N]: Lift member updates out of consumer and
share TargetAssignmentBuilder (#22510)
6fa3711c3e9 is described below
commit 6fa3711c3e94bcc50104f08996db4346060412a6
Author: Sean Quah <[email protected]>
AuthorDate: Mon Jun 29 12:55:43 2026 +0100
KAFKA-20292 [3/N]: Lift member updates out of consumer and share
TargetAssignmentBuilder (#22510)
Remove the member update logic from the consumer and share group
TargetAssignmentBuilder. Instead, use the new
UpdatedMembersAndTargetAssignmentView class to provide updated views of
members and the target assignment to the TargetAssignmentBuilder.
Reviewers: Dongnuo Lyu <[email protected]>, David Jacot
<[email protected]>
---
.../coordinator/group/GroupMetadataManager.java | 39 +-
.../group/modern/TargetAssignmentBuilder.java | 76 ----
.../group/modern/TargetAssignmentBuilderTest.java | 438 +--------------------
.../assignor/TargetAssignmentBuilderBenchmark.java | 13 +-
4 files changed, 34 insertions(+), 532 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 08a0abfda04..78b42edf754 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
@@ -166,6 +166,7 @@ import
org.apache.kafka.coordinator.group.streams.topics.ConfiguredSubtopology;
import org.apache.kafka.coordinator.group.streams.topics.ConfiguredTopology;
import org.apache.kafka.coordinator.group.streams.topics.InternalTopicManager;
import
org.apache.kafka.coordinator.group.streams.topics.TopicConfigurationException;
+import
org.apache.kafka.coordinator.group.util.UpdatedMembersAndTargetAssignmentView;
import org.apache.kafka.server.authorizer.AuthorizableRequestContext;
import org.apache.kafka.server.authorizer.Authorizer;
import org.apache.kafka.server.share.persister.DeleteShareGroupStateParameters;
@@ -4145,24 +4146,23 @@ public class GroupMetadataManager {
updatedMember
).orElse(defaultConsumerGroupAssignor.name());
try {
+ UpdatedMembersAndTargetAssignmentView<ConsumerGroupMember,
Assignment> updatedMembersAndTargetAssignment =
+ new UpdatedMembersAndTargetAssignmentView<>(
+ group.members(),
+ group.staticMembers(),
+ group.targetAssignment()
+ );
+
updatedMembersAndTargetAssignment.addOrUpdateMember(updatedMember.memberId(),
updatedMember.instanceId(), updatedMember);
+
TargetAssignmentBuilder.ConsumerTargetAssignmentBuilder
assignmentResultBuilder =
new
TargetAssignmentBuilder.ConsumerTargetAssignmentBuilder(group.groupId(),
groupEpoch, consumerGroupAssignors.get(preferredServerAssignor))
.withTime(time)
- .withMembers(group.members())
- .withStaticMembers(group.staticMembers())
+ .withMembers(updatedMembersAndTargetAssignment.members())
.withSubscriptionType(subscriptionType)
- .withTargetAssignment(group.targetAssignment())
+
.withTargetAssignment(updatedMembersAndTargetAssignment.targetAssignment())
.withInvertedTargetAssignment(group.invertedTargetAssignment())
.withMetadataImage(metadataImage)
-
.withResolvedRegularExpressions(group.resolvedRegularExpressions())
- .addOrUpdateMember(updatedMember.memberId(),
updatedMember);
-
- // 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.
- String previousMemberId =
group.staticMemberId(updatedMember.instanceId());
- if (previousMemberId != null &&
!updatedMember.memberId().equals(previousMemberId)) {
- assignmentResultBuilder.removeMember(previousMemberId);
- }
+
.withResolvedRegularExpressions(group.resolvedRegularExpressions());
long startTimeMs = time.milliseconds();
TargetAssignmentBuilder.TargetAssignmentResult assignmentResult =
@@ -4229,16 +4229,23 @@ public class GroupMetadataManager {
stripInitValue(shareGroupStatePartitionMetadata.get(group.groupId()).initializedTopics())
:
Map.of();
+ UpdatedMembersAndTargetAssignmentView<ShareGroupMember,
Assignment> updatedMembersAndTargetAssignment =
+ new UpdatedMembersAndTargetAssignmentView<>(
+ group.members(),
+ Map.of(),
+ group.targetAssignment()
+ );
+
updatedMembersAndTargetAssignment.addOrUpdateMember(updatedMember.memberId(),
updatedMember.instanceId(), updatedMember);
+
TargetAssignmentBuilder.ShareTargetAssignmentBuilder
assignmentResultBuilder =
new
TargetAssignmentBuilder.ShareTargetAssignmentBuilder(group.groupId(),
groupEpoch, shareGroupAssignor)
.withTime(time)
- .withMembers(group.members())
+ .withMembers(updatedMembersAndTargetAssignment.members())
.withSubscriptionType(subscriptionType)
- .withTargetAssignment(group.targetAssignment())
+
.withTargetAssignment(updatedMembersAndTargetAssignment.targetAssignment())
.withTopicAssignablePartitionsMap(initializedTopicPartitions)
.withInvertedTargetAssignment(group.invertedTargetAssignment())
- .withMetadataImage(metadataImage)
- .addOrUpdateMember(updatedMember.memberId(),
updatedMember);
+ .withMetadataImage(metadataImage);
long startTimeMs = time.milliseconds();
TargetAssignmentBuilder.TargetAssignmentResult assignmentResult =
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/modern/TargetAssignmentBuilder.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/modern/TargetAssignmentBuilder.java
index 86c64e5c22d..6937b95eaf6 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/modern/TargetAssignmentBuilder.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/modern/TargetAssignmentBuilder.java
@@ -287,17 +287,6 @@ public abstract class TargetAssignmentBuilder<T extends
ModernGroupMember, U ext
*/
private CoordinatorMetadataImage metadataImage =
CoordinatorMetadataImage.EMPTY;
- /**
- * The members which have been updated or deleted. Deleted members
- * are signaled by a null value.
- */
- private final Map<String, T> updatedMembers = new HashMap<>();
-
- /**
- * The static members in the group.
- */
- private Map<String, String> staticMembers = new HashMap<>();
-
/**
* Topic partition assignable map.
*/
@@ -344,19 +333,6 @@ public abstract class TargetAssignmentBuilder<T extends
ModernGroupMember, U ext
return self();
}
- /**
- * Adds all the existing static members.
- *
- * @param staticMembers The existing static members in the consumer
group.
- * @return This object.
- */
- public U withStaticMembers(
- Map<String, String> staticMembers
- ) {
- this.staticMembers = staticMembers;
- return self();
- }
-
/**
* Adds the subscription type in use.
*
@@ -416,35 +392,6 @@ public abstract class TargetAssignmentBuilder<T extends
ModernGroupMember, U ext
return self();
}
- /**
- * 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 U addOrUpdateMember(
- String memberId,
- T member
- ) {
- this.updatedMembers.put(memberId, member);
- return self();
- }
-
- /**
- * 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 U removeMember(
- String memberId
- ) {
- return addOrUpdateMember(memberId, null);
- }
-
/**
* Builds the new target assignment.
*
@@ -465,29 +412,6 @@ public abstract class TargetAssignmentBuilder<T extends
ModernGroupMember, U ext
))
);
- // Update the member spec if updated or deleted members.
- updatedMembers.forEach((memberId, updatedMemberOrNull) -> {
- if (updatedMemberOrNull == null) {
- memberSpecs.remove(memberId);
- } else {
- Assignment assignment =
targetAssignment.getOrDefault(memberId, Assignment.EMPTY);
-
- // A new static member joins and needs to replace an existing
departed one.
- if (updatedMemberOrNull.instanceId() != null) {
- String previousMemberId =
staticMembers.get(updatedMemberOrNull.instanceId());
- if (previousMemberId != null &&
!previousMemberId.equals(memberId)) {
- assignment =
targetAssignment.getOrDefault(previousMemberId, Assignment.EMPTY);
- }
- }
-
- memberSpecs.put(memberId, newMemberSubscriptionAndAssignment(
- updatedMemberOrNull,
- assignment,
- topicResolver
- ));
- }
- });
-
// Compute the assignment.
GroupAssignment newGroupAssignment = assignor.assign(
new GroupSpecImpl(
diff --git
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/modern/TargetAssignmentBuilderTest.java
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/modern/TargetAssignmentBuilderTest.java
index c575b621e4d..b13ea51f441 100644
---
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/modern/TargetAssignmentBuilderTest.java
+++
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/modern/TargetAssignmentBuilderTest.java
@@ -38,7 +38,6 @@ import java.util.Map;
import java.util.Optional;
import java.util.Set;
-import static
org.apache.kafka.coordinator.group.Assertions.assertRecordsEquals;
import static
org.apache.kafka.coordinator.group.Assertions.assertUnorderedRecordsEquals;
import static
org.apache.kafka.coordinator.group.AssignmentTestUtil.mkAssignment;
import static
org.apache.kafka.coordinator.group.AssignmentTestUtil.mkTopicAssignment;
@@ -60,10 +59,8 @@ public class TargetAssignmentBuilderTest {
private final long assignmentTimestamp;
private final PartitionAssignor assignor =
mock(PartitionAssignor.class);
private final Map<String, ConsumerGroupMember> members = new
HashMap<>();
- private final Map<String, ConsumerGroupMember> updatedMembers = new
HashMap<>();
private final Map<String, Assignment> targetAssignment = new
HashMap<>();
private final Map<String, MemberAssignment> memberAssignments = new
HashMap<>();
- private final Map<String, String> staticMembers = new HashMap<>();
private final Map<String, ResolvedRegularExpression>
resolvedRegularExpressions = new HashMap<>();
private MetadataImageBuilder metadataImageBuilder = new
MetadataImageBuilder();
@@ -82,7 +79,7 @@ public class TargetAssignmentBuilderTest {
List<String> subscriptions,
Map<Uuid, Set<Integer>> targetPartitions
) {
- addGroupMember(memberId, null, subscriptions, "",
targetPartitions);
+ addGroupMember(memberId, subscriptions, "", targetPartitions);
}
public void addGroupMember(
@@ -90,34 +87,10 @@ public class TargetAssignmentBuilderTest {
List<String> subscriptions,
String subscribedRegex,
Map<Uuid, Set<Integer>> targetPartitions
- ) {
- addGroupMember(memberId, null, subscriptions, subscribedRegex,
targetPartitions);
- }
-
- public void addGroupMember(
- String memberId,
- String instanceId,
- List<String> subscriptions,
- Map<Uuid, Set<Integer>> targetPartitions
- ) {
- addGroupMember(memberId, instanceId, subscriptions, "",
targetPartitions);
- }
-
- public void addGroupMember(
- String memberId,
- String instanceId,
- List<String> subscriptions,
- String subscribedRegex,
- Map<Uuid, Set<Integer>> targetPartitions
) {
ConsumerGroupMember.Builder memberBuilder = new
ConsumerGroupMember.Builder(memberId)
.setSubscribedTopicNames(subscriptions)
.setSubscribedTopicRegex(subscribedRegex);
-
- if (instanceId != null) {
- memberBuilder.setInstanceId(instanceId);
- staticMembers.put(instanceId, memberId);
- }
members.put(memberId, memberBuilder.build());
targetAssignment.put(memberId, new Assignment(targetPartitions));
}
@@ -132,44 +105,6 @@ public class TargetAssignmentBuilderTest {
return topicId;
}
- public void updateMemberSubscription(
- String memberId,
- List<String> subscriptions
- ) {
- updateMemberSubscription(
- memberId,
- subscriptions,
- Optional.empty(),
- Optional.empty()
- );
- }
-
- public void updateMemberSubscription(
- String memberId,
- List<String> subscriptions,
- Optional<String> instanceId,
- Optional<String> rackId
- ) {
- ConsumerGroupMember existingMember = members.get(memberId);
- ConsumerGroupMember.Builder builder;
- if (existingMember != null) {
- builder = new ConsumerGroupMember.Builder(existingMember);
- } else {
- builder = new ConsumerGroupMember.Builder(memberId);
- }
- updatedMembers.put(memberId, builder
- .setSubscribedTopicNames(subscriptions)
- .maybeUpdateInstanceId(instanceId)
- .maybeUpdateRackId(rackId)
- .build());
- }
-
- public void removeMemberSubscription(
- String memberId
- ) {
- this.updatedMembers.put(memberId, null);
- }
-
public void prepareMemberAssignment(
String memberId,
Map<Uuid, Set<Integer>> assignment
@@ -219,10 +154,9 @@ public class TargetAssignmentBuilderTest {
public TargetAssignmentBuilder.TargetAssignmentResult build() {
CoordinatorMetadataImage coordinatorMetadataImage = new
KRaftCoordinatorMetadataImage(metadataImageBuilder.build());
TopicIds.TopicResolver topicResolver = new
TopicIds.CachedTopicResolver(coordinatorMetadataImage);
+
// Prepare expected member specs.
Map<String, MemberSubscriptionAndAssignmentImpl>
memberSubscriptions = new HashMap<>();
-
- // All the existing members are prepared.
members.forEach((memberId, member) ->
memberSubscriptions.put(memberId,
newMemberSubscriptionAndAssignment(
member,
@@ -231,30 +165,6 @@ public class TargetAssignmentBuilderTest {
))
);
- // All the updated are added and all the deleted
- // members are removed.
- updatedMembers.forEach((memberId, updatedMemberOrNull) -> {
- if (updatedMemberOrNull == null) {
- memberSubscriptions.remove(memberId);
- } else {
- Assignment assignment =
targetAssignment.getOrDefault(memberId, Assignment.EMPTY);
-
- // A new static member joins and needs to replace an
existing departed one.
- if (updatedMemberOrNull.instanceId() != null) {
- String previousMemberId =
staticMembers.get(updatedMemberOrNull.instanceId());
- if (previousMemberId != null &&
!previousMemberId.equals(memberId)) {
- assignment =
targetAssignment.getOrDefault(previousMemberId, Assignment.EMPTY);
- }
- }
-
- memberSubscriptions.put(memberId,
newMemberSubscriptionAndAssignment(
- updatedMemberOrNull,
- assignment,
- topicResolver
- ));
- }
- });
-
// Prepare the expected subscription topic metadata.
SubscribedTopicDescriberImpl subscribedTopicMetadata = new
SubscribedTopicDescriberImpl(coordinatorMetadataImage);
SubscriptionType subscriptionType = HOMOGENEOUS;
@@ -280,22 +190,12 @@ public class TargetAssignmentBuilderTest {
new
TargetAssignmentBuilder.ConsumerTargetAssignmentBuilder(groupId, groupEpoch,
assignor)
.withTime(new MockTime(0, assignmentTimestamp,
assignmentTimestamp))
.withMembers(members)
- .withStaticMembers(staticMembers)
.withSubscriptionType(subscriptionType)
.withTargetAssignment(targetAssignment)
.withInvertedTargetAssignment(invertedTargetAssignment)
.withMetadataImage(coordinatorMetadataImage)
.withResolvedRegularExpressions(resolvedRegularExpressions);
- // 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.
TargetAssignmentBuilder.TargetAssignmentResult result =
builder.build();
@@ -446,183 +346,6 @@ public class TargetAssignmentBuilderTest {
assertEquals(expectedAssignment, result.targetAssignment());
}
- @Test
- public void testNewMember() {
- TargetAssignmentBuilderTestContext context = new
TargetAssignmentBuilderTestContext(
- "my-group",
- 20,
- 12345L
- );
-
- Uuid fooTopicId = context.addTopicMetadata("foo", 6);
- Uuid barTopicId = context.addTopicMetadata("bar", 6);
-
- context.addGroupMember("member-1", Arrays.asList("foo", "bar", "zar"),
mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2, 3),
- mkTopicAssignment(barTopicId, 1, 2, 3)
- ));
-
- context.addGroupMember("member-2", Arrays.asList("foo", "bar", "zar"),
mkAssignment(
- mkTopicAssignment(fooTopicId, 4, 5, 6),
- mkTopicAssignment(barTopicId, 4, 5, 6)
- ));
-
- context.updateMemberSubscription("member-3", Arrays.asList("foo",
"bar", "zar"));
-
- context.prepareMemberAssignment("member-1", mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2),
- mkTopicAssignment(barTopicId, 1, 2)
- ));
-
- context.prepareMemberAssignment("member-2", mkAssignment(
- mkTopicAssignment(fooTopicId, 3, 4),
- mkTopicAssignment(barTopicId, 3, 4)
- ));
-
- context.prepareMemberAssignment("member-3", mkAssignment(
- mkTopicAssignment(fooTopicId, 5, 6),
- mkTopicAssignment(barTopicId, 5, 6)
- ));
-
- TargetAssignmentBuilder.TargetAssignmentResult result =
context.build();
-
- assertUnorderedRecordsEquals(
- List.of(
- List.of(
- newConsumerGroupTargetAssignmentRecord("my-group",
"member-1", mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2),
- mkTopicAssignment(barTopicId, 1, 2)
- )),
- newConsumerGroupTargetAssignmentRecord("my-group",
"member-2", mkAssignment(
- mkTopicAssignment(fooTopicId, 3, 4),
- mkTopicAssignment(barTopicId, 3, 4)
- )),
- newConsumerGroupTargetAssignmentRecord("my-group",
"member-3", mkAssignment(
- mkTopicAssignment(fooTopicId, 5, 6),
- mkTopicAssignment(barTopicId, 5, 6)
- ))
- ),
- List.of(
- newConsumerGroupTargetAssignmentMetadataRecord(
- "my-group",
- 20,
- 12345L
- )
- )
- ),
- result.records()
- );
-
- Map<String, MemberAssignment> expectedAssignment = new HashMap<>();
- expectedAssignment.put("member-1", new
MemberAssignmentImpl(mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2),
- mkTopicAssignment(barTopicId, 1, 2)
- )));
- expectedAssignment.put("member-2", new
MemberAssignmentImpl(mkAssignment(
- mkTopicAssignment(fooTopicId, 3, 4),
- mkTopicAssignment(barTopicId, 3, 4)
- )));
- expectedAssignment.put("member-3", new
MemberAssignmentImpl(mkAssignment(
- mkTopicAssignment(fooTopicId, 5, 6),
- mkTopicAssignment(barTopicId, 5, 6)
- )));
-
- assertEquals(expectedAssignment, result.targetAssignment());
- }
-
- @Test
- public void testUpdateMember() {
- TargetAssignmentBuilderTestContext context = new
TargetAssignmentBuilderTestContext(
- "my-group",
- 20,
- 12345L
- );
-
- Uuid fooTopicId = context.addTopicMetadata("foo", 6);
- Uuid barTopicId = context.addTopicMetadata("bar", 6);
-
- context.addGroupMember("member-1", Arrays.asList("foo", "bar", "zar"),
mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2, 3),
- mkTopicAssignment(barTopicId, 1, 2)
- ));
-
- context.addGroupMember("member-2", Arrays.asList("foo", "bar", "zar"),
mkAssignment(
- mkTopicAssignment(fooTopicId, 4, 5, 6),
- mkTopicAssignment(barTopicId, 3, 4)
- ));
-
- context.addGroupMember("member-3", Arrays.asList("bar", "zar"),
mkAssignment(
- mkTopicAssignment(barTopicId, 5, 6)
- ));
-
- context.updateMemberSubscription(
- "member-3",
- Arrays.asList("foo", "bar", "zar"),
- Optional.of("instance-id-3"),
- Optional.of("rack-0")
- );
-
- context.prepareMemberAssignment("member-1", mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2),
- mkTopicAssignment(barTopicId, 1, 2)
- ));
-
- context.prepareMemberAssignment("member-2", mkAssignment(
- mkTopicAssignment(fooTopicId, 3, 4),
- mkTopicAssignment(barTopicId, 3, 4)
- ));
-
- context.prepareMemberAssignment("member-3", mkAssignment(
- mkTopicAssignment(fooTopicId, 5, 6),
- mkTopicAssignment(barTopicId, 5, 6)
- ));
-
- TargetAssignmentBuilder.TargetAssignmentResult result =
context.build();
-
- assertUnorderedRecordsEquals(
- List.of(
- List.of(
- newConsumerGroupTargetAssignmentRecord("my-group",
"member-1", mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2),
- mkTopicAssignment(barTopicId, 1, 2)
- )),
- newConsumerGroupTargetAssignmentRecord("my-group",
"member-2", mkAssignment(
- mkTopicAssignment(fooTopicId, 3, 4),
- mkTopicAssignment(barTopicId, 3, 4)
- )),
- newConsumerGroupTargetAssignmentRecord("my-group",
"member-3", mkAssignment(
- mkTopicAssignment(fooTopicId, 5, 6),
- mkTopicAssignment(barTopicId, 5, 6)
- ))
- ),
- List.of(
- newConsumerGroupTargetAssignmentMetadataRecord(
- "my-group",
- 20,
- 12345L
- )
- )
- ),
- result.records()
- );
-
- Map<String, MemberAssignment> expectedAssignment = new HashMap<>();
- expectedAssignment.put("member-1", new
MemberAssignmentImpl(mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2),
- mkTopicAssignment(barTopicId, 1, 2)
- )));
- expectedAssignment.put("member-2", new
MemberAssignmentImpl(mkAssignment(
- mkTopicAssignment(fooTopicId, 3, 4),
- mkTopicAssignment(barTopicId, 3, 4)
- )));
- expectedAssignment.put("member-3", new
MemberAssignmentImpl(mkAssignment(
- mkTopicAssignment(fooTopicId, 5, 6),
- mkTopicAssignment(barTopicId, 5, 6)
- )));
-
- assertEquals(expectedAssignment, result.targetAssignment());
- }
-
@Test
public void testPartialAssignmentUpdate() {
TargetAssignmentBuilderTestContext context = new
TargetAssignmentBuilderTestContext(
@@ -707,163 +430,6 @@ public class TargetAssignmentBuilderTest {
assertEquals(expectedAssignment, result.targetAssignment());
}
- @Test
- public void testDeleteMember() {
- TargetAssignmentBuilderTestContext context = new
TargetAssignmentBuilderTestContext(
- "my-group",
- 20,
- 12345L
- );
-
- Uuid fooTopicId = context.addTopicMetadata("foo", 6);
- Uuid barTopicId = context.addTopicMetadata("bar", 6);
-
- context.addGroupMember("member-1", Arrays.asList("foo", "bar", "zar"),
mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2),
- mkTopicAssignment(barTopicId, 1, 2)
- ));
-
- context.addGroupMember("member-2", Arrays.asList("foo", "bar", "zar"),
mkAssignment(
- mkTopicAssignment(fooTopicId, 3, 4),
- mkTopicAssignment(barTopicId, 3, 4)
- ));
-
- context.addGroupMember("member-3", Arrays.asList("foo", "bar", "zar"),
mkAssignment(
- mkTopicAssignment(fooTopicId, 5, 6),
- mkTopicAssignment(barTopicId, 5, 6)
- ));
-
- context.removeMemberSubscription("member-3");
-
- context.prepareMemberAssignment("member-1", mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2, 3),
- mkTopicAssignment(barTopicId, 1, 2, 3)
- ));
-
- context.prepareMemberAssignment("member-2", mkAssignment(
- mkTopicAssignment(fooTopicId, 4, 5, 6),
- mkTopicAssignment(barTopicId, 4, 5, 6)
- ));
-
- TargetAssignmentBuilder.TargetAssignmentResult result =
context.build();
-
- assertUnorderedRecordsEquals(
- List.of(
- List.of(
- newConsumerGroupTargetAssignmentRecord("my-group",
"member-1", mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2, 3),
- mkTopicAssignment(barTopicId, 1, 2, 3)
- )),
- newConsumerGroupTargetAssignmentRecord("my-group",
"member-2", mkAssignment(
- mkTopicAssignment(fooTopicId, 4, 5, 6),
- mkTopicAssignment(barTopicId, 4, 5, 6)
- ))
- ),
- List.of(
- newConsumerGroupTargetAssignmentMetadataRecord(
- "my-group",
- 20,
- 12345L
- )
- )
- ),
- result.records()
- );
-
- Map<String, MemberAssignment> expectedAssignment = new HashMap<>();
- expectedAssignment.put("member-1", new
MemberAssignmentImpl(mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2, 3),
- mkTopicAssignment(barTopicId, 1, 2, 3)
- )));
- expectedAssignment.put("member-2", new
MemberAssignmentImpl(mkAssignment(
- mkTopicAssignment(fooTopicId, 4, 5, 6),
- mkTopicAssignment(barTopicId, 4, 5, 6)
- )));
-
- assertEquals(expectedAssignment, result.targetAssignment());
- }
-
- @Test
- public void testReplaceStaticMember() {
- TargetAssignmentBuilderTestContext context = new
TargetAssignmentBuilderTestContext(
- "my-group",
- 20,
- 12345L
- );
-
- Uuid fooTopicId = context.addTopicMetadata("foo", 6);
- Uuid barTopicId = context.addTopicMetadata("bar", 6);
-
- context.addGroupMember("member-1", "instance-member-1",
Arrays.asList("foo", "bar", "zar"), mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2),
- mkTopicAssignment(barTopicId, 1, 2)
- ));
-
- context.addGroupMember("member-2", "instance-member-2",
Arrays.asList("foo", "bar", "zar"), mkAssignment(
- mkTopicAssignment(fooTopicId, 3, 4),
- mkTopicAssignment(barTopicId, 3, 4)
- ));
-
- context.addGroupMember("member-3", "instance-member-3",
Arrays.asList("foo", "bar", "zar"), mkAssignment(
- mkTopicAssignment(fooTopicId, 5, 6),
- mkTopicAssignment(barTopicId, 5, 6)
- ));
-
- // Static member 3 leaves
- context.removeMemberSubscription("member-3");
-
- // Another static member joins with the same instance id as the
departed one
- context.updateMemberSubscription("member-3-a", Arrays.asList("foo",
"bar", "zar"), Optional.of("instance-member-3"), Optional.empty());
-
- context.prepareMemberAssignment("member-1", mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2),
- mkTopicAssignment(barTopicId, 1, 2)
- ));
-
- context.prepareMemberAssignment("member-2", mkAssignment(
- mkTopicAssignment(fooTopicId, 3, 4),
- mkTopicAssignment(barTopicId, 3, 4)
- ));
-
- context.prepareMemberAssignment("member-3-a", mkAssignment(
- mkTopicAssignment(fooTopicId, 5, 6),
- mkTopicAssignment(barTopicId, 5, 6)
- ));
-
- TargetAssignmentBuilder.TargetAssignmentResult result =
context.build();
-
- assertRecordsEquals(
- List.of(
- newConsumerGroupTargetAssignmentRecord("my-group",
"member-3-a", mkAssignment(
- mkTopicAssignment(fooTopicId, 5, 6),
- mkTopicAssignment(barTopicId, 5, 6)
- )),
- newConsumerGroupTargetAssignmentMetadataRecord(
- "my-group",
- 20,
- 12345L
- )
- ),
- result.records()
- );
-
- Map<String, MemberAssignment> expectedAssignment = new HashMap<>();
- expectedAssignment.put("member-1", new
MemberAssignmentImpl(mkAssignment(
- mkTopicAssignment(fooTopicId, 1, 2),
- mkTopicAssignment(barTopicId, 1, 2)
- )));
- expectedAssignment.put("member-2", new
MemberAssignmentImpl(mkAssignment(
- mkTopicAssignment(fooTopicId, 3, 4),
- mkTopicAssignment(barTopicId, 3, 4)
- )));
-
- expectedAssignment.put("member-3-a", new
MemberAssignmentImpl(mkAssignment(
- mkTopicAssignment(fooTopicId, 5, 6),
- mkTopicAssignment(barTopicId, 5, 6)
- )));
-
- assertEquals(expectedAssignment, result.targetAssignment());
- }
@Test
public void testRegularExpressions() {
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 b7c1f49bf91..2859b2f3961 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
@@ -31,6 +31,7 @@ import
org.apache.kafka.coordinator.group.modern.SubscribedTopicDescriberImpl;
import org.apache.kafka.coordinator.group.modern.TargetAssignmentBuilder;
import org.apache.kafka.coordinator.group.modern.TopicIds;
import org.apache.kafka.coordinator.group.modern.consumer.ConsumerGroupMember;
+import
org.apache.kafka.coordinator.group.util.UpdatedMembersAndTargetAssignmentView;
import org.openjdk.jmh.annotations.Benchmark;
import org.openjdk.jmh.annotations.BenchmarkMode;
@@ -81,6 +82,8 @@ public class TargetAssignmentBuilderBenchmark {
private PartitionAssignor partitionAssignor;
+ private UpdatedMembersAndTargetAssignmentView<ConsumerGroupMember,
Assignment> updatedMembersAndTargetAssignment;
+
private TargetAssignmentBuilder.ConsumerTargetAssignmentBuilder
targetAssignmentBuilder;
/** The number of homogeneous subgroups to create for the heterogeneous
subscription case. */
@@ -113,14 +116,16 @@ public class TargetAssignmentBuilderBenchmark {
.setSubscribedTopicNames(allTopicNames)
.build();
+ updatedMembersAndTargetAssignment = new
UpdatedMembersAndTargetAssignmentView<>(members, Map.of(),
existingTargetAssignment);
+
updatedMembersAndTargetAssignment.addOrUpdateMember(newMember.memberId(),
newMember.instanceId(), newMember);
+
targetAssignmentBuilder = new
TargetAssignmentBuilder.ConsumerTargetAssignmentBuilder(GROUP_ID, GROUP_EPOCH,
partitionAssignor)
.withTime(Time.SYSTEM)
- .withMembers(members)
+ .withMembers(updatedMembersAndTargetAssignment.members())
.withSubscriptionType(subscriptionType)
- .withTargetAssignment(existingTargetAssignment)
+
.withTargetAssignment(updatedMembersAndTargetAssignment.targetAssignment())
.withInvertedTargetAssignment(invertedTargetAssignment)
- .withMetadataImage(metadataImage)
- .addOrUpdateMember(newMember.memberId(), newMember);
+ .withMetadataImage(metadataImage);
}
private void setupTopics() {