diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupAssignment.java b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/GroupAssignment.java similarity index 82% rename from group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupAssignment.java rename to group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/GroupAssignment.java index c4ac5803b35d0..19a84e3471809 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupAssignment.java +++ b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/GroupAssignment.java @@ -14,7 +14,10 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.coordinator.group.streams.assignor; +package org.apache.kafka.coordinator.group.api.streams.assignor; + +import org.apache.kafka.common.annotation.InterfaceAudience; +import org.apache.kafka.common.annotation.InterfaceStability; import java.util.Map; import java.util.Objects; @@ -24,6 +27,8 @@ * * @param members The member assignments keyed by member ID. */ +@InterfaceAudience.Public +@InterfaceStability.Evolving public record GroupAssignment(Map members) { public GroupAssignment { diff --git a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/GroupSpec.java b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/GroupSpec.java new file mode 100644 index 0000000000000..486b2f8e17999 --- /dev/null +++ b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/GroupSpec.java @@ -0,0 +1,58 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.coordinator.group.api.streams.assignor; + +import org.apache.kafka.common.annotation.InterfaceAudience; +import org.apache.kafka.common.annotation.InterfaceStability; + +import java.util.Collection; +import java.util.Map; + +/** + * The group metadata specifications required to compute the target assignment. + */ +@InterfaceAudience.Public +@InterfaceStability.Evolving +public interface GroupSpec { + + /** + * @return The member Ids of all members in the group. + */ + Collection memberIds(); + + /** + * Gets the static metadata for a given member. + * + * @param memberId The member Id. + * @return The static member metadata. + */ + MemberAssignmentMetadata memberMetadata(String memberId); + + /** + * Gets the current assignment state for a given member. + * + * @param memberId The member Id. + * @return The current member assignment state. + */ + MemberAssignmentState memberAssignmentState(String memberId); + + /** + * @return Any configurations passed to the assignor. + */ + Map configs(); + +} diff --git a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignment.java b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignment.java new file mode 100644 index 0000000000000..8dc7c63c1bcce --- /dev/null +++ b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignment.java @@ -0,0 +1,45 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.coordinator.group.api.streams.assignor; + +import org.apache.kafka.common.annotation.InterfaceAudience; +import org.apache.kafka.common.annotation.InterfaceStability; + +import java.util.Map; +import java.util.Set; + +/** + * The task assignment for a streams group member. + * + *

Only active and standby tasks are assigned by the {@link TaskAssignor}. Warm-up tasks are + * not assigned by the assignor; they are decided during reconciliation. + */ +@InterfaceAudience.Public +@InterfaceStability.Evolving +public interface MemberAssignment { + + /** + * @return The active tasks assigned to this member keyed by subtopology Id. + */ + Map> activeTasks(); + + /** + * @return The standby tasks assigned to this member keyed by subtopology Id. + */ + Map> standbyTasks(); + +} diff --git a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentMetadata.java b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentMetadata.java new file mode 100644 index 0000000000000..ef725075d6dc0 --- /dev/null +++ b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentMetadata.java @@ -0,0 +1,56 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.coordinator.group.api.streams.assignor; + +import org.apache.kafka.common.annotation.InterfaceAudience; +import org.apache.kafka.common.annotation.InterfaceStability; + +import java.util.Map; +import java.util.Optional; + +/** + * Interface representing the static metadata for a streams group member. + * + *

The metadata contains the per-member information that does not change during the assignment + * computation. The member's current task assignment state is exposed separately through + * {@link MemberAssignmentState}. + */ +@InterfaceAudience.Public +@InterfaceStability.Evolving +public interface MemberAssignmentMetadata { + + /** + * @return The instance ID if provided. + */ + Optional instanceId(); + + /** + * @return The rack ID if provided. + */ + Optional rackId(); + + /** + * @return The process ID. + */ + String processId(); + + /** + * @return The client tags for a rack-aware assignment. + */ + Map clientTags(); + +} diff --git a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentState.java b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentState.java new file mode 100644 index 0000000000000..c21c4db119bdc --- /dev/null +++ b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentState.java @@ -0,0 +1,73 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.coordinator.group.api.streams.assignor; + +import org.apache.kafka.common.annotation.InterfaceAudience; +import org.apache.kafka.common.annotation.InterfaceStability; + +import java.util.Map; +import java.util.Set; + +/** + * Interface representing the current assignment state for a streams group member, used by the + * {@link TaskAssignor} to compute a new, sticky target assignment. + * + *

The active and standby tasks are the member's current target assignment (the last + * assignment committed for the member). Using the committed target as the stickiness baseline keeps + * in-flight task moves converging rather than reverting them. Warm-up tasks are not assigned by the + * assignor (they are decided during reconciliation); the tasks the member is currently + * warming up are exposed here so that the assignor can take them into account. As a result, a task + * that is being moved to this member can appear both as an active or standby task (from the target) + * and as a warm-up task (currently in progress). + * + *

All accessors are keyed by subtopology ID. The task-set accessors ({@link #activeTasks()}, + * {@link #standbyTasks()}, {@link #warmupTasks()}) map each subtopology ID to its set of + * partitions, while {@link #taskOffsets()} and {@link #taskEndOffsets()} map each subtopology ID + * to a map from partition to offset. + */ +@InterfaceAudience.Public +@InterfaceStability.Evolving +public interface MemberAssignmentState { + + /** + * @return The member's current target active tasks keyed by subtopology Id. + */ + Map> activeTasks(); + + /** + * @return The member's current target standby tasks keyed by subtopology Id. + */ + Map> standbyTasks(); + + /** + * @return The tasks the member is currently warming up, keyed by subtopology Id. + */ + Map> warmupTasks(); + + /** + * @return The last received cumulative task offsets of assigned tasks or dormant tasks. + * The outer map is keyed by subtopology ID and the inner map by partition. + */ + Map> taskOffsets(); + + /** + * @return The last received cumulative task end offsets of assigned tasks or dormant tasks. + * The outer map is keyed by subtopology ID and the inner map by partition. + */ + Map> taskEndOffsets(); + +} diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TaskAssignor.java b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TaskAssignor.java similarity index 77% rename from group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TaskAssignor.java rename to group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TaskAssignor.java index 73bff075ca648..f6532b8c84026 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TaskAssignor.java +++ b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TaskAssignor.java @@ -14,15 +14,20 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.coordinator.group.streams.assignor; +package org.apache.kafka.coordinator.group.api.streams.assignor; + +import org.apache.kafka.common.annotation.InterfaceAudience; +import org.apache.kafka.common.annotation.InterfaceStability; /** * Server side task assignor used by streams groups. */ +@InterfaceAudience.Public +@InterfaceStability.Evolving public interface TaskAssignor { /** - * Unique name for this assignor. + * Unique name for this assignor. Used in configuration to select this assignor. */ String name(); @@ -33,7 +38,7 @@ public interface TaskAssignor { * @param topologyDescriber The task metadata describer. * @return The new assignment for the group. * - * @throws TaskAssignorException For empty groups + * @throws TaskAssignorException If the assignment cannot be computed. */ GroupAssignment assign( GroupSpec groupSpec, diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TaskAssignorException.java b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TaskAssignorException.java similarity index 73% rename from group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TaskAssignorException.java rename to group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TaskAssignorException.java index 2fda6a9e9ec62..872ab7825c0e4 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TaskAssignorException.java +++ b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TaskAssignorException.java @@ -14,13 +14,19 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.coordinator.group.streams.assignor; +package org.apache.kafka.coordinator.group.api.streams.assignor; +import org.apache.kafka.common.annotation.InterfaceAudience; +import org.apache.kafka.common.annotation.InterfaceStability; import org.apache.kafka.common.errors.ApiException; /** - * Exception thrown by {@link TaskAssignor#assign(GroupSpec, TopologyDescriber)}}. The exception is only used internally. + * Exception thrown by {@link TaskAssignor#assign(GroupSpec, TopologyDescriber)} when the group's tasks + * cannot be assigned. Custom {@link TaskAssignor} implementations should throw this exception to signal + * an assignment failure. */ +@InterfaceAudience.Public +@InterfaceStability.Evolving public class TaskAssignorException extends ApiException { public TaskAssignorException(String message) { diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TopologyDescriber.java b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TopologyDescriber.java similarity index 85% rename from group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TopologyDescriber.java rename to group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TopologyDescriber.java index d6f7a3ab579c6..326de17755011 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TopologyDescriber.java +++ b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TopologyDescriber.java @@ -14,7 +14,10 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.coordinator.group.streams.assignor; +package org.apache.kafka.coordinator.group.api.streams.assignor; + +import org.apache.kafka.common.annotation.InterfaceAudience; +import org.apache.kafka.common.annotation.InterfaceStability; import java.util.List; import java.util.NoSuchElementException; @@ -22,12 +25,14 @@ /** * The topology describer is used by the {@link TaskAssignor} to get topic and task metadata of the group's topology. */ +@InterfaceAudience.Public +@InterfaceStability.Evolving public interface TopologyDescriber { /** - * Map of topic names to topic metadata. + * The IDs of all subtopologies in the group's topology. * - * @return The list of subtopologies IDs. + * @return The list of subtopology IDs. */ List subtopologies(); diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpec.java b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/package-info.java similarity index 65% rename from group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpec.java rename to group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/package-info.java index 1a8e7edc01c63..b8c31e7e537e4 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpec.java +++ b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/package-info.java @@ -14,23 +14,8 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.kafka.coordinator.group.streams.assignor; - -import java.util.Map; /** - * The group metadata specifications required to compute the target assignment. + * Provides the public API for broker-side custom task assignors used by Kafka Streams groups. */ -public interface GroupSpec { - - /** - * @return Member metadata keyed by member Id. - */ - Map members(); - - /** - * @return Any configurations passed to the assignor. - */ - Map assignmentConfigs(); - -} +package org.apache.kafka.coordinator.group.api.streams.assignor; 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 f2442c70463f9..43ef8e9033179 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 @@ -92,6 +92,8 @@ import org.apache.kafka.coordinator.group.api.assignor.PartitionAssignorException; import org.apache.kafka.coordinator.group.api.assignor.ShareGroupPartitionAssignor; import org.apache.kafka.coordinator.group.api.assignor.SubscriptionType; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignorException; import org.apache.kafka.coordinator.group.assignor.SimpleAssignor; import org.apache.kafka.coordinator.group.classic.ClassicGroup; import org.apache.kafka.coordinator.group.classic.ClassicGroupMember; @@ -162,8 +164,6 @@ import org.apache.kafka.coordinator.group.streams.TasksTuple; import org.apache.kafka.coordinator.group.streams.TasksTupleWithEpochs; import org.apache.kafka.coordinator.group.streams.assignor.StickyTaskAssignor; -import org.apache.kafka.coordinator.group.streams.assignor.TaskAssignor; -import org.apache.kafka.coordinator.group.streams.assignor.TaskAssignorException; 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; diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsets.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsets.java index ed2434a79bffa..82ebe9d096169 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsets.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsets.java @@ -18,7 +18,6 @@ import org.apache.kafka.common.errors.InvalidRequestException; import org.apache.kafka.common.message.StreamsGroupHeartbeatRequestData; -import org.apache.kafka.coordinator.group.streams.assignor.TaskId; import java.util.HashMap; import java.util.List; @@ -32,10 +31,10 @@ * constantly changing"); they are held in memory on the group coordinator and re-reported by the member on the * task-offset interval. * - * @param taskOffsets Cumulative changelog offsets per task. - * @param taskEndOffsets Cumulative changelog end-offsets per task. + * @param taskOffsets Cumulative changelog offsets, keyed by subtopology ID then partition. + * @param taskEndOffsets Cumulative changelog end-offsets, keyed by subtopology ID then partition. */ -public record MemberTaskOffsets(Map taskOffsets, Map taskEndOffsets) { +public record MemberTaskOffsets(Map> taskOffsets, Map> taskEndOffsets) { public static final MemberTaskOffsets EMPTY = new MemberTaskOffsets(Map.of(), Map.of()); @@ -52,19 +51,19 @@ public MemberTaskOffsets update( final List reportedTaskEndOffsets ) { return new MemberTaskOffsets( - reportedTaskOffsets == null ? taskOffsets : toTaskIdMap(reportedTaskOffsets), - reportedTaskEndOffsets == null ? taskEndOffsets : toTaskIdMap(reportedTaskEndOffsets) + reportedTaskOffsets == null ? taskOffsets : toNestedMap(reportedTaskOffsets), + reportedTaskEndOffsets == null ? taskEndOffsets : toNestedMap(reportedTaskEndOffsets) ); } - private static Map toTaskIdMap(final List taskOffsets) { - final Map result = new HashMap<>(taskOffsets.size()); + private static Map> toNestedMap(final List taskOffsets) { + final Map> result = new HashMap<>(); for (final StreamsGroupHeartbeatRequestData.TaskOffset taskOffset : taskOffsets) { - final TaskId taskId = new TaskId(taskOffset.subtopologyId(), taskOffset.partition()); + final Map byPartition = result.computeIfAbsent(taskOffset.subtopologyId(), k -> new HashMap<>()); // The reported values come straight from the client heartbeat and the protocol does not enforce uniqueness // of (subtopologyId, partition). Reject a duplicate with a clear client error rather than silently picking // one of the values. - if (result.putIfAbsent(taskId, taskOffset.offset()) != null) { + if (byPartition.putIfAbsent(taskOffset.partition(), taskOffset.offset()) != null) { throw new InvalidRequestException("Task offsets contain a duplicate entry for subtopology " + taskOffset.subtopologyId() + " and partition " + taskOffset.partition() + "."); } 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 6ce2632b6f986..23bab5f3e3a42 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 @@ -19,12 +19,13 @@ import org.apache.kafka.common.utils.Time; import org.apache.kafka.coordinator.common.runtime.CoordinatorMetadataImage; import org.apache.kafka.coordinator.common.runtime.CoordinatorRecord; -import org.apache.kafka.coordinator.group.streams.assignor.AssignmentMemberSpec; -import org.apache.kafka.coordinator.group.streams.assignor.GroupAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignorException; import org.apache.kafka.coordinator.group.streams.assignor.GroupSpecImpl; -import org.apache.kafka.coordinator.group.streams.assignor.MemberAssignment; -import org.apache.kafka.coordinator.group.streams.assignor.TaskAssignor; -import org.apache.kafka.coordinator.group.streams.assignor.TaskAssignorException; +import org.apache.kafka.coordinator.group.streams.assignor.MemberAssignmentImpl; +import org.apache.kafka.coordinator.group.streams.assignor.MemberMetadataAndAssignmentImpl; import org.apache.kafka.coordinator.group.streams.topics.ConfiguredTopology; import java.util.ArrayList; @@ -116,17 +117,19 @@ public TargetAssignmentBuilder( this.assignmentConfigs = Objects.requireNonNull(assignmentConfigs); } - static AssignmentMemberSpec createAssignmentMemberSpec( + static MemberMetadataAndAssignmentImpl createMemberMetadataAndAssignment( StreamsGroupMember member, TasksTuple targetAssignment, MemberTaskOffsets taskOffsets ) { - return new AssignmentMemberSpec( + return new MemberMetadataAndAssignmentImpl( member.instanceId(), member.rackId(), targetAssignment.activeTasks(), targetAssignment.standbyTasks(), - targetAssignment.warmupTasks(), + // Warm-up tasks are decided during reconciliation and recorded in the member's current + // assignment, not in the target assignment, so they are sourced from there. + member.assignedTasks().warmupTasks(), member.processId(), member.clientTags(), taskOffsets.taskOffsets(), @@ -217,10 +220,10 @@ public TargetAssignmentBuilder withTopology( * @throws TaskAssignorException if the target assignment cannot be computed. */ public TargetAssignmentResult build() throws TaskAssignorException { - Map memberSpecs = new HashMap<>(); + Map memberMetadataMap = new HashMap<>(); - // Prepare the member spec for all members. - members.forEach((memberId, member) -> memberSpecs.put(memberId, createAssignmentMemberSpec( + // Prepare the member metadata for all members. + members.forEach((memberId, member) -> memberMetadataMap.put(memberId, createMemberMetadataAndAssignment( member, targetAssignment.getOrDefault(memberId, org.apache.kafka.coordinator.group.streams.TasksTuple.EMPTY), taskOffsets.getOrDefault(memberId, MemberTaskOffsets.EMPTY) @@ -234,14 +237,14 @@ public TargetAssignmentResult build() throws TaskAssignorException { } newGroupAssignment = assignor.assign( new GroupSpecImpl( - Collections.unmodifiableMap(memberSpecs), + Collections.unmodifiableMap(memberMetadataMap), assignmentConfigs ), new TopologyMetadata(metadataImage, topology.subtopologies().get()) ); } else { newGroupAssignment = new GroupAssignment( - memberSpecs.keySet().stream().collect(Collectors.toMap(x -> x, x -> MemberAssignment.empty()))); + memberMetadataMap.keySet().stream().collect(Collectors.toMap(x -> x, x -> (MemberAssignment) MemberAssignmentImpl.empty()))); } // Compute delta from previous to new target assignment and create the @@ -249,7 +252,7 @@ public TargetAssignmentResult build() throws TaskAssignorException { List records = new ArrayList<>(); Map newTargetAssignment = new HashMap<>(); - memberSpecs.keySet().forEach(memberId -> { + memberMetadataMap.keySet().forEach(memberId -> { org.apache.kafka.coordinator.group.streams.TasksTuple oldMemberAssignment = targetAssignment.get(memberId); org.apache.kafka.coordinator.group.streams.TasksTuple newMemberAssignment = newMemberAssignment(newGroupAssignment, memberId); @@ -294,7 +297,8 @@ private TasksTuple newMemberAssignment( return new TasksTuple( newMemberAssignment.activeTasks(), newMemberAssignment.standbyTasks(), - newMemberAssignment.warmupTasks() + // Warm-up tasks are not assigned by the assignor; they are decided during reconciliation. + Map.of() ); } else { return TasksTuple.EMPTY; diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TopologyMetadata.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TopologyMetadata.java index e07a3f1eca42f..2286392bb003c 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TopologyMetadata.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TopologyMetadata.java @@ -17,7 +17,7 @@ package org.apache.kafka.coordinator.group.streams; import org.apache.kafka.coordinator.common.runtime.CoordinatorMetadataImage; -import org.apache.kafka.coordinator.group.streams.assignor.TopologyDescriber; +import org.apache.kafka.coordinator.group.api.streams.assignor.TopologyDescriber; import org.apache.kafka.coordinator.group.streams.topics.ConfiguredSubtopology; import java.util.Collections; @@ -27,7 +27,7 @@ import java.util.SortedMap; /** - * The topology metadata class is used by the {@link org.apache.kafka.coordinator.group.streams.assignor.TaskAssignor} to get topic and + * The topology metadata class is used by the {@link org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignor} to get topic and * partition metadata for the topology that the streams group using. * * @param metadataImage The metadata image diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/AssignmentMemberSpec.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/AssignmentMemberSpec.java deleted file mode 100644 index fa41b5a8a3b4b..0000000000000 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/AssignmentMemberSpec.java +++ /dev/null @@ -1,60 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.apache.kafka.coordinator.group.streams.assignor; - -import java.util.Collections; -import java.util.Map; -import java.util.Objects; -import java.util.Optional; -import java.util.Set; - -/** - * The assignment specification for a streams group member. - * - * @param instanceId The instance ID if provided. - * @param rackId The rack ID if provided. - * @param activeTasks Current target active tasks - * @param standbyTasks Current target standby tasks - * @param warmupTasks Current target warm-up tasks - * @param processId The process ID. - * @param clientTags The client tags for a rack-aware assignment. - * @param taskOffsets The last received cumulative task offsets of assigned tasks or dormant tasks. - */ -public record AssignmentMemberSpec(Optional instanceId, - Optional rackId, - Map> activeTasks, - Map> standbyTasks, - Map> warmupTasks, - String processId, - Map clientTags, - Map taskOffsets, - Map taskEndOffsets -) { - - public AssignmentMemberSpec { - Objects.requireNonNull(instanceId); - Objects.requireNonNull(rackId); - activeTasks = Collections.unmodifiableMap(Objects.requireNonNull(activeTasks)); - standbyTasks = Collections.unmodifiableMap(Objects.requireNonNull(standbyTasks)); - warmupTasks = Collections.unmodifiableMap(Objects.requireNonNull(warmupTasks)); - Objects.requireNonNull(processId); - clientTags = Collections.unmodifiableMap(Objects.requireNonNull(clientTags)); - taskOffsets = Collections.unmodifiableMap(Objects.requireNonNull(taskOffsets)); - taskEndOffsets = Collections.unmodifiableMap(Objects.requireNonNull(taskEndOffsets)); - } - -} diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImpl.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImpl.java index 7479fcae0fc2d..c271d2c72186b 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImpl.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImpl.java @@ -16,21 +16,52 @@ */ package org.apache.kafka.coordinator.group.streams.assignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupSpec; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignmentMetadata; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignmentState; + +import java.util.Collection; import java.util.Map; import java.util.Objects; /** * The assignment specification for a streams group. * - * @param members The member metadata keyed by member ID. - * @param assignmentConfigs Any configurations passed to the assignor. + * @param members The member metadata keyed by member Id. Each value provides both the + * {@link MemberAssignmentMetadata} and the {@link MemberAssignmentState} for the member. + * @param configs Any configurations passed to the assignor. */ -public record GroupSpecImpl(Map members, - Map assignmentConfigs) implements GroupSpec { +public record GroupSpecImpl( + Map members, + Map configs +) implements GroupSpec { public GroupSpecImpl { Objects.requireNonNull(members); - Objects.requireNonNull(assignmentConfigs); + Objects.requireNonNull(configs); + } + + @Override + public Collection memberIds() { + return members.keySet(); + } + + @Override + public MemberAssignmentMetadata memberMetadata(String memberId) { + return requireMember(memberId); + } + + @Override + public MemberAssignmentState memberAssignmentState(String memberId) { + return requireMember(memberId); + } + + private MemberMetadataAndAssignmentImpl requireMember(String memberId) { + MemberMetadataAndAssignmentImpl member = members.get(memberId); + if (member == null) { + throw new IllegalArgumentException("Member Id " + memberId + " not found."); + } + return member; } } diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberAssignment.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberAssignmentImpl.java similarity index 64% rename from group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberAssignment.java rename to group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberAssignmentImpl.java index 2902e647382f7..be81adb1f3fb7 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberAssignment.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberAssignmentImpl.java @@ -16,6 +16,9 @@ */ package org.apache.kafka.coordinator.group.streams.assignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignment; + +import java.util.HashMap; import java.util.Map; import java.util.Objects; import java.util.Set; @@ -23,21 +26,21 @@ /** * The task assignment for a streams group member. * - * @param activeTasks The active tasks assigned to this member keyed by subtopologyId. - * @param standbyTasks The standby tasks assigned to this member keyed by subtopologyId. - * @param warmupTasks The warm-up tasks assigned to this member keyed by subtopologyId. + * @param activeTasks The active tasks assigned to this member keyed by subtopology Id. + * @param standbyTasks The standby tasks assigned to this member keyed by subtopology Id. + * The maps are not made immutable, since the server-side assignors rely on + * being able to mutate them while building new assignments. */ -public record MemberAssignment(Map> activeTasks, - Map> standbyTasks, - Map> warmupTasks) { +public record MemberAssignmentImpl(Map> activeTasks, + Map> standbyTasks) implements MemberAssignment { - public MemberAssignment { + public MemberAssignmentImpl { Objects.requireNonNull(activeTasks); Objects.requireNonNull(standbyTasks); - Objects.requireNonNull(warmupTasks); } - public static MemberAssignment empty() { - return new MemberAssignment(Map.of(), Map.of(), Map.of()); + public static MemberAssignmentImpl empty() { + return new MemberAssignmentImpl(new HashMap<>(), new HashMap<>()); } + } diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberMetadataAndAssignmentImpl.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberMetadataAndAssignmentImpl.java new file mode 100644 index 0000000000000..0d3a031eb1dd8 --- /dev/null +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberMetadataAndAssignmentImpl.java @@ -0,0 +1,65 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.coordinator.group.streams.assignor; + +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignmentMetadata; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignmentState; + +import java.util.Collections; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.Set; + +/** + * Implementation of both the {@link MemberAssignmentMetadata} and the {@link MemberAssignmentState} + * interfaces for a streams group member. + * + * @param instanceId The instance ID if provided. + * @param rackId The rack ID if provided. + * @param activeTasks Current active tasks. + * @param standbyTasks Current standby tasks. + * @param warmupTasks Current warm-up tasks. + * @param processId The process ID. + * @param clientTags The client tags for a rack-aware assignment. + * @param taskOffsets The last received cumulative task offsets of assigned tasks or dormant tasks. + * @param taskEndOffsets The last received cumulative task end offsets of assigned tasks or dormant tasks. + */ +public record MemberMetadataAndAssignmentImpl(Optional instanceId, + Optional rackId, + Map> activeTasks, + Map> standbyTasks, + Map> warmupTasks, + String processId, + Map clientTags, + Map> taskOffsets, + Map> taskEndOffsets +) implements MemberAssignmentMetadata, MemberAssignmentState { + + public MemberMetadataAndAssignmentImpl { + Objects.requireNonNull(instanceId); + Objects.requireNonNull(rackId); + activeTasks = Collections.unmodifiableMap(Objects.requireNonNull(activeTasks)); + standbyTasks = Collections.unmodifiableMap(Objects.requireNonNull(standbyTasks)); + warmupTasks = Collections.unmodifiableMap(Objects.requireNonNull(warmupTasks)); + Objects.requireNonNull(processId); + clientTags = Collections.unmodifiableMap(Objects.requireNonNull(clientTags)); + taskOffsets = Collections.unmodifiableMap(Objects.requireNonNull(taskOffsets)); + taskEndOffsets = Collections.unmodifiableMap(Objects.requireNonNull(taskEndOffsets)); + } + +} diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MockAssignor.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MockAssignor.java index 93827d29b1c75..284d59826474e 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MockAssignor.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MockAssignor.java @@ -16,6 +16,13 @@ */ package org.apache.kafka.coordinator.group.streams.assignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupSpec; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignorException; +import org.apache.kafka.coordinator.group.api.streams.assignor.TopologyDescriber; + import java.util.Comparator; import java.util.HashMap; import java.util.HashSet; @@ -56,13 +63,10 @@ public GroupAssignment assign( } // Copy existing assignment and fill temporary data structures - for (Map.Entry memberEntry : groupSpec.members().entrySet()) { - final String memberId = memberEntry.getKey(); - final AssignmentMemberSpec memberSpec = memberEntry.getValue(); - - Map> activeTasks = new HashMap<>(memberSpec.activeTasks()); + for (final String memberId : groupSpec.memberIds()) { + Map> activeTasks = new HashMap<>(groupSpec.memberAssignmentState(memberId).activeTasks()); - newTargetAssignment.put(memberId, new MemberAssignment(activeTasks, new HashMap<>(), new HashMap<>())); + newTargetAssignment.put(memberId, new MemberAssignmentImpl(activeTasks, new HashMap<>())); for (Map.Entry> entry : activeTasks.entrySet()) { final String subtopologyId = entry.getKey(); final Set taskIds = entry.getValue(); diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/ProcessState.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/ProcessState.java index 330a0cc0da3a7..b06ff236f49a4 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/ProcessState.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/ProcessState.java @@ -16,6 +16,8 @@ */ package org.apache.kafka.coordinator.group.streams.assignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignorException; + import java.util.AbstractMap; import java.util.HashMap; import java.util.HashSet; diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignor.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignor.java index 4828f0238b266..03bedad1f37c4 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignor.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignor.java @@ -17,10 +17,20 @@ package org.apache.kafka.coordinator.group.streams.assignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupSpec; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignmentMetadata; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignmentState; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignorException; +import org.apache.kafka.coordinator.group.api.streams.assignor.TopologyDescriber; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.ArrayList; +import java.util.Collection; import java.util.Comparator; import java.util.HashMap; import java.util.HashSet; @@ -66,7 +76,7 @@ private GroupAssignment doAssign(final GroupSpec groupSpec, final TopologyDescri assignStandby(statefulTasks); } - return buildGroupAssignment(groupSpec.members().keySet()); + return buildGroupAssignment(groupSpec.memberIds()); } private LinkedList taskIds(final TopologyDescriber topologyDescriber, final boolean isActive) { @@ -85,8 +95,8 @@ private LinkedList taskIds(final TopologyDescriber topologyDescriber, fi private void initialize(final GroupSpec groupSpec, final TopologyDescriber topologyDescriber) { localState = new LocalState(); localState.numStandbyReplicas = - groupSpec.assignmentConfigs().isEmpty() ? 0 - : Integer.parseInt(groupSpec.assignmentConfigs().get("num.standby.replicas")); + groupSpec.configs().isEmpty() ? 0 + : Integer.parseInt(groupSpec.configs().get("num.standby.replicas")); // Helpers for computing active tasks per member, and tasks per member localState.totalActiveTasks = 0; @@ -98,25 +108,25 @@ private void initialize(final GroupSpec groupSpec, final TopologyDescriber topol if (topologyDescriber.isStateful(subtopology)) localState.totalTasks += numberOfPartitions * localState.numStandbyReplicas; } - localState.totalMembersWithActiveTaskCapacity = groupSpec.members().size(); - localState.totalMembersWithTaskCapacity = groupSpec.members().size(); + localState.totalMembersWithActiveTaskCapacity = groupSpec.memberIds().size(); + localState.totalMembersWithTaskCapacity = groupSpec.memberIds().size(); localState.activeTasksPerMember = computeTasksPerMember(localState.totalActiveTasks, localState.totalMembersWithActiveTaskCapacity); localState.totalTasksPerMember = computeTasksPerMember(localState.totalTasks, localState.totalMembersWithTaskCapacity); localState.processIdToState = new HashMap<>(localState.totalMembersWithActiveTaskCapacity); localState.activeTaskToPrevMember = new HashMap<>(localState.totalActiveTasks); localState.standbyTaskToPrevMember = new HashMap<>(localState.numStandbyReplicas > 0 ? (localState.totalTasks - localState.totalActiveTasks) / localState.numStandbyReplicas : 0); - for (final Map.Entry memberEntry : groupSpec.members().entrySet()) { - final String memberId = memberEntry.getKey(); - final String processId = memberEntry.getValue().processId(); + for (final String memberId : groupSpec.memberIds()) { + final MemberAssignmentMetadata memberMetadata = groupSpec.memberMetadata(memberId); + final MemberAssignmentState memberAssignmentState = groupSpec.memberAssignmentState(memberId); + final String processId = memberMetadata.processId(); final Member member = new Member(processId, memberId); - final AssignmentMemberSpec memberSpec = memberEntry.getValue(); localState.processIdToState.putIfAbsent(processId, new ProcessState(processId)); localState.processIdToState.get(processId).addMember(memberId); // prev active tasks - for (final Map.Entry> entry : memberSpec.activeTasks().entrySet()) { + for (final Map.Entry> entry : memberAssignmentState.activeTasks().entrySet()) { final Set partitionNoSet = entry.getValue(); for (final int partitionNo : partitionNoSet) { localState.activeTaskToPrevMember.put(new TaskId(entry.getKey(), partitionNo), member); @@ -124,7 +134,7 @@ private void initialize(final GroupSpec groupSpec, final TopologyDescriber topol } // prev standby tasks - for (final Map.Entry> entry : memberSpec.standbyTasks().entrySet()) { + for (final Map.Entry> entry : memberAssignmentState.standbyTasks().entrySet()) { final Set partitionNoSet = entry.getValue(); for (final int partitionNo : partitionNoSet) { final TaskId taskId = new TaskId(entry.getKey(), partitionNo); @@ -135,7 +145,7 @@ private void initialize(final GroupSpec groupSpec, final TopologyDescriber topol } } - private GroupAssignment buildGroupAssignment(final Set members) { + private GroupAssignment buildGroupAssignment(final Collection members) { final Map memberAssignments = new HashMap<>(); final Map> activeTasksAssignments = localState.processIdToState.entrySet().stream() @@ -162,7 +172,7 @@ private GroupAssignment buildGroupAssignment(final Set members) { if (standbyTasksAssignments.containsKey(memberId)) { standByTasks.putAll(toCompactedTaskIds(standbyTasksAssignments.get(memberId))); } - memberAssignments.put(memberId, new MemberAssignment(activeTasks, standByTasks, new HashMap<>())); + memberAssignments.put(memberId, new MemberAssignmentImpl(activeTasks, standByTasks)); } return new GroupAssignment(memberAssignments); diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TaskId.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TaskId.java index 979c042df00c5..fe1a14b105b62 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TaskId.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/TaskId.java @@ -20,7 +20,7 @@ import java.util.Objects; /** - * The identifier for a task + * The identifier for a task, consisting of the subtopology ID and the partition. * * @param subtopologyId The unique identifier of the subtopology. * @param partition The partition of the input topics this task is processing. 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 04841537e421d..904da6b023b1b 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 @@ -97,6 +97,8 @@ import org.apache.kafka.coordinator.group.api.assignor.GroupAssignment; import org.apache.kafka.coordinator.group.api.assignor.GroupSpec; import org.apache.kafka.coordinator.group.api.assignor.PartitionAssignorException; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignorException; import org.apache.kafka.coordinator.group.classic.ClassicGroup; import org.apache.kafka.coordinator.group.classic.ClassicGroupMember; import org.apache.kafka.coordinator.group.classic.ClassicGroupState; @@ -148,8 +150,6 @@ import org.apache.kafka.coordinator.group.streams.TaskAssignmentTestUtil.TaskRole; import org.apache.kafka.coordinator.group.streams.TasksTuple; import org.apache.kafka.coordinator.group.streams.TasksTupleWithEpochs; -import org.apache.kafka.coordinator.group.streams.assignor.TaskAssignor; -import org.apache.kafka.coordinator.group.streams.assignor.TaskAssignorException; import org.apache.kafka.image.MetadataDelta; import org.apache.kafka.image.MetadataImage; import org.apache.kafka.image.MetadataProvenance; @@ -18991,8 +18991,8 @@ fooTopicName, computeTopicHash(fooTopicName, metadataImage) StreamsGroup group = context.groupMetadataManager.streamsGroup(groupId); assertEquals( new org.apache.kafka.coordinator.group.streams.MemberTaskOffsets( - Map.of(new org.apache.kafka.coordinator.group.streams.assignor.TaskId(subtopology1, 0), 10L), - Map.of(new org.apache.kafka.coordinator.group.streams.assignor.TaskId(subtopology1, 0), 20L) + Map.of(subtopology1, Map.of(0, 10L)), + Map.of(subtopology1, Map.of(0, 20L)) ), group.taskOffsets(memberId) ); @@ -19017,8 +19017,8 @@ fooTopicName, computeTopicHash(fooTopicName, metadataImage) assertEquals(List.of(), result.records()); assertEquals( new org.apache.kafka.coordinator.group.streams.MemberTaskOffsets( - Map.of(new org.apache.kafka.coordinator.group.streams.assignor.TaskId(subtopology1, 0), 12L), - Map.of(new org.apache.kafka.coordinator.group.streams.assignor.TaskId(subtopology1, 0), 20L) + Map.of(subtopology1, Map.of(0, 12L)), + Map.of(subtopology1, Map.of(0, 20L)) ), group.taskOffsets(memberId) ); @@ -19592,8 +19592,9 @@ public void testStreamsGroupMemberJoiningWithStaleTopology() { ) .build(); - assignor.prepareGroupAssignment(new org.apache.kafka.coordinator.group.streams.assignor.GroupAssignment(Map.of( - memberId, org.apache.kafka.coordinator.group.streams.assignor.MemberAssignment.empty() + assignor.prepareGroupAssignment(new org.apache.kafka.coordinator.group.api.streams.assignor.GroupAssignment(Map.of( + memberId, (org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignment) + org.apache.kafka.coordinator.group.streams.assignor.MemberAssignmentImpl.empty() ))); // Member joins the streams group. diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTestContext.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTestContext.java index 0b8ea094c2da8..c7a5ff4812715 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTestContext.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTestContext.java @@ -60,6 +60,7 @@ import org.apache.kafka.coordinator.common.runtime.MockCoordinatorTimer; import org.apache.kafka.coordinator.group.api.assignor.ConsumerGroupPartitionAssignor; import org.apache.kafka.coordinator.group.api.assignor.ShareGroupPartitionAssignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignor; import org.apache.kafka.coordinator.group.classic.ClassicGroup; import org.apache.kafka.coordinator.group.generated.ConsumerGroupCurrentMemberAssignmentKey; import org.apache.kafka.coordinator.group.generated.ConsumerGroupCurrentMemberAssignmentValue; @@ -114,7 +115,6 @@ import org.apache.kafka.coordinator.group.streams.StreamsGroupHeartbeatResult; import org.apache.kafka.coordinator.group.streams.StreamsGroupMember; import org.apache.kafka.coordinator.group.streams.TasksTupleWithEpochs; -import org.apache.kafka.coordinator.group.streams.assignor.TaskAssignor; import org.apache.kafka.coordinator.group.streams.topics.InternalTopicManager; import org.apache.kafka.server.authorizer.Authorizer; import org.apache.kafka.server.common.ApiMessageAndVersion; diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsetsTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsetsTest.java index 941c5c3cc97a0..143e3079690b6 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsetsTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsetsTest.java @@ -18,7 +18,6 @@ import org.apache.kafka.common.errors.InvalidRequestException; import org.apache.kafka.common.message.StreamsGroupHeartbeatRequestData; -import org.apache.kafka.coordinator.group.streams.assignor.TaskId; import org.junit.jupiter.api.Test; @@ -47,11 +46,11 @@ public void shouldConvertHeartbeatRequestTaskOffsets() { ); assertEquals( - Map.of(new TaskId("sub-1", 0), 10L, new TaskId("sub-1", 1), 20L), + Map.of("sub-1", Map.of(0, 10L, 1, 20L)), result.taskOffsets() ); assertEquals( - Map.of(new TaskId("sub-1", 0), 15L, new TaskId("sub-1", 1), 25L), + Map.of("sub-1", Map.of(0, 15L, 1, 25L)), result.taskEndOffsets() ); } @@ -76,8 +75,8 @@ public void shouldRejectDuplicateTaskEntries() { @Test public void shouldRetainBothMapsWhenBothListsAreNull() { MemberTaskOffsets previous = new MemberTaskOffsets( - Map.of(new TaskId("sub-1", 0), 10L), - Map.of(new TaskId("sub-1", 0), 15L) + Map.of("sub-1", Map.of(0, 10L)), + Map.of("sub-1", Map.of(0, 15L)) ); assertEquals(previous, previous.update(null, null)); @@ -87,28 +86,28 @@ public void shouldRetainBothMapsWhenBothListsAreNull() { @Test public void shouldUpdateTaskOffsetsAndRetainTaskEndOffsetsWhenEndOffsetsNull() { MemberTaskOffsets previous = new MemberTaskOffsets( - Map.of(new TaskId("sub-1", 0), 10L), - Map.of(new TaskId("sub-1", 0), 15L) + Map.of("sub-1", Map.of(0, 10L)), + Map.of("sub-1", Map.of(0, 15L)) ); MemberTaskOffsets result = previous.update(List.of(taskOffset("sub-1", 0, 12L)), null); - assertEquals(Map.of(new TaskId("sub-1", 0), 12L), result.taskOffsets()); + assertEquals(Map.of("sub-1", Map.of(0, 12L)), result.taskOffsets()); // The end-offsets were not reported, so the previously reported values are retained. - assertEquals(Map.of(new TaskId("sub-1", 0), 15L), result.taskEndOffsets()); + assertEquals(Map.of("sub-1", Map.of(0, 15L)), result.taskEndOffsets()); } @Test public void shouldUpdateTaskEndOffsetsAndRetainTaskOffsetsWhenOffsetsNull() { MemberTaskOffsets previous = new MemberTaskOffsets( - Map.of(new TaskId("sub-1", 0), 10L), - Map.of(new TaskId("sub-1", 0), 15L) + Map.of("sub-1", Map.of(0, 10L)), + Map.of("sub-1", Map.of(0, 15L)) ); MemberTaskOffsets result = previous.update(null, List.of(taskOffset("sub-1", 0, 18L))); // The offsets were not reported, so the previously reported values are retained. - assertEquals(Map.of(new TaskId("sub-1", 0), 10L), result.taskOffsets()); - assertEquals(Map.of(new TaskId("sub-1", 0), 18L), result.taskEndOffsets()); + assertEquals(Map.of("sub-1", Map.of(0, 10L)), result.taskOffsets()); + assertEquals(Map.of("sub-1", Map.of(0, 18L)), result.taskEndOffsets()); } } diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MockTaskAssignor.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MockTaskAssignor.java index 2b1f1fb551183..6a4a8d19e1423 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MockTaskAssignor.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MockTaskAssignor.java @@ -16,12 +16,13 @@ */ package org.apache.kafka.coordinator.group.streams; -import org.apache.kafka.coordinator.group.streams.assignor.GroupAssignment; -import org.apache.kafka.coordinator.group.streams.assignor.GroupSpec; -import org.apache.kafka.coordinator.group.streams.assignor.MemberAssignment; -import org.apache.kafka.coordinator.group.streams.assignor.TaskAssignor; -import org.apache.kafka.coordinator.group.streams.assignor.TaskAssignorException; -import org.apache.kafka.coordinator.group.streams.assignor.TopologyDescriber; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupSpec; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignorException; +import org.apache.kafka.coordinator.group.api.streams.assignor.TopologyDescriber; +import org.apache.kafka.coordinator.group.streams.assignor.MemberAssignmentImpl; import java.util.Map; import java.util.Map.Entry; @@ -48,8 +49,8 @@ public void prepareGroupAssignment(Map memberAssignments) { Entry::getKey, entry -> { TasksTuple tasksTuple = entry.getValue(); - return new MemberAssignment( - tasksTuple.activeTasks(), tasksTuple.standbyTasks(), tasksTuple.warmupTasks()); + return (MemberAssignment) new MemberAssignmentImpl( + tasksTuple.activeTasks(), tasksTuple.standbyTasks()); }))); } @@ -70,7 +71,7 @@ public String toString() { @Override public GroupAssignment assign(final GroupSpec groupSpec, final TopologyDescriber topologyDescriber) throws TaskAssignorException { - assignmentConfigs = groupSpec.assignmentConfigs(); + assignmentConfigs = groupSpec.configs(); return preparedGroupAssignment; } } diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/StreamsGroupTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/StreamsGroupTest.java index 8d79f8e1dedcf..9f81bed979626 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/StreamsGroupTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/StreamsGroupTest.java @@ -48,7 +48,6 @@ import org.apache.kafka.coordinator.group.generated.StreamsGroupTopologyValue; import org.apache.kafka.coordinator.group.streams.StreamsGroup.StreamsGroupState; import org.apache.kafka.coordinator.group.streams.TaskAssignmentTestUtil.TaskRole; -import org.apache.kafka.coordinator.group.streams.assignor.TaskId; import org.apache.kafka.coordinator.group.streams.topics.ConfiguredTopology; import org.apache.kafka.image.MetadataImage; import org.apache.kafka.timeline.SnapshotRegistry; @@ -120,8 +119,8 @@ public void testUpdateAndRetrieveTaskOffsets() { assertEquals(Map.of(), streamsGroup.taskOffsets()); MemberTaskOffsets offsets = new MemberTaskOffsets( - Map.of(new TaskId("sub-1", 0), 10L), - Map.of(new TaskId("sub-1", 0), 20L) + Map.of("sub-1", Map.of(0, 10L)), + Map.of("sub-1", Map.of(0, 20L)) ); streamsGroup.updateTaskOffsets("member-id", offsets); @@ -130,8 +129,8 @@ public void testUpdateAndRetrieveTaskOffsets() { // A new report replaces the previous one. MemberTaskOffsets newerOffsets = new MemberTaskOffsets( - Map.of(new TaskId("sub-1", 0), 15L), - Map.of(new TaskId("sub-1", 0), 25L) + Map.of("sub-1", Map.of(0, 15L)), + Map.of("sub-1", Map.of(0, 25L)) ); streamsGroup.updateTaskOffsets("member-id", newerOffsets); assertEquals(newerOffsets, streamsGroup.taskOffsets("member-id")); @@ -142,8 +141,8 @@ public void testRemoveMemberClearsTaskOffsets() { StreamsGroup streamsGroup = createStreamsGroup("foo"); streamsGroup.updateMember(new StreamsGroupMember.Builder("member-id").build()); streamsGroup.updateTaskOffsets("member-id", new MemberTaskOffsets( - Map.of(new TaskId("sub-1", 0), 10L), - Map.of(new TaskId("sub-1", 0), 20L) + Map.of("sub-1", Map.of(0, 10L)), + Map.of("sub-1", Map.of(0, 20L)) )); streamsGroup.removeMember("member-id"); 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 0edb840870c19..fbe2835fca0b8 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 @@ -22,14 +22,14 @@ import org.apache.kafka.coordinator.common.runtime.CoordinatorRecord; import org.apache.kafka.coordinator.common.runtime.KRaftCoordinatorMetadataImage; import org.apache.kafka.coordinator.common.runtime.MetadataImageBuilder; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignor; import org.apache.kafka.coordinator.group.generated.StreamsGroupMemberMetadataValue; import org.apache.kafka.coordinator.group.streams.TaskAssignmentTestUtil.TaskRole; -import org.apache.kafka.coordinator.group.streams.assignor.AssignmentMemberSpec; -import org.apache.kafka.coordinator.group.streams.assignor.GroupAssignment; import org.apache.kafka.coordinator.group.streams.assignor.GroupSpecImpl; -import org.apache.kafka.coordinator.group.streams.assignor.MemberAssignment; -import org.apache.kafka.coordinator.group.streams.assignor.TaskAssignor; -import org.apache.kafka.coordinator.group.streams.assignor.TaskId; +import org.apache.kafka.coordinator.group.streams.assignor.MemberAssignmentImpl; +import org.apache.kafka.coordinator.group.streams.assignor.MemberMetadataAndAssignmentImpl; import org.apache.kafka.coordinator.group.streams.topics.ConfiguredSubtopology; import org.apache.kafka.coordinator.group.streams.topics.ConfiguredTopology; @@ -50,7 +50,7 @@ import static org.apache.kafka.coordinator.group.Assertions.assertUnorderedRecordsEquals; import static org.apache.kafka.coordinator.group.streams.StreamsCoordinatorRecordHelpers.newStreamsGroupTargetAssignmentMetadataRecord; import static org.apache.kafka.coordinator.group.streams.StreamsCoordinatorRecordHelpers.newStreamsGroupTargetAssignmentRecord; -import static org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder.createAssignmentMemberSpec; +import static org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder.createMemberMetadataAndAssignment; import static org.apache.kafka.coordinator.group.streams.TaskAssignmentTestUtil.mkTasks; import static org.apache.kafka.coordinator.group.streams.TaskAssignmentTestUtil.mkTasksTuple; import static org.junit.jupiter.api.Assertions.assertEquals; @@ -87,8 +87,10 @@ public void testBuildEmptyAssignmentWhenTopologyNotReady() { } @ParameterizedTest - @EnumSource(TaskRole.class) - public void testCreateAssignmentMemberSpec(TaskRole taskRole) { + // Warm-up tasks are sourced from the member's current assignment (set during reconciliation), + // not from the target assignment, so this test only varies the target's active and standby tasks. + @EnumSource(value = TaskRole.class, names = {"ACTIVE", "STANDBY"}) + public void testCreateMemberMetadataAndAssignment(TaskRole taskRole) { String fooSubtopologyId = Uuid.randomUuid().toString(); String barSubtopologyId = Uuid.randomUuid().toString(); @@ -98,6 +100,7 @@ public void testCreateAssignmentMemberSpec(TaskRole taskRole) { .setInstanceId("instanceId") .setProcessId("processId") .setClientTags(clientTags) + .setAssignedTasks(TasksTupleWithEpochs.EMPTY) .build(); TasksTuple assignment = mkTasksTuple(taskRole, @@ -105,27 +108,27 @@ public void testCreateAssignmentMemberSpec(TaskRole taskRole) { mkTasks(barSubtopologyId, 1, 2, 3) ); - AssignmentMemberSpec assignmentMemberSpec = createAssignmentMemberSpec( + MemberMetadataAndAssignmentImpl memberMetadata = createMemberMetadataAndAssignment( member, assignment, MemberTaskOffsets.EMPTY ); - assertEquals(new AssignmentMemberSpec( + assertEquals(new MemberMetadataAndAssignmentImpl( Optional.of("instanceId"), Optional.of("rackId"), assignment.activeTasks(), assignment.standbyTasks(), - assignment.warmupTasks(), + Map.of(), "processId", clientTags, Map.of(), Map.of() - ), assignmentMemberSpec); + ), memberMetadata); } @Test - public void testCreateAssignmentMemberSpecPopulatesTaskOffsets() { + public void testCreateMemberMetadataAndAssignmentPopulatesTaskOffsets() { String fooSubtopologyId = Uuid.randomUuid().toString(); StreamsGroupMember member = new StreamsGroupMember.Builder("member-id") @@ -133,19 +136,44 @@ public void testCreateAssignmentMemberSpecPopulatesTaskOffsets() { .setInstanceId("instanceId") .setProcessId("processId") .setClientTags(Map.of()) + .setAssignedTasks(TasksTupleWithEpochs.EMPTY) .build(); - Map taskOffsets = Map.of(new TaskId(fooSubtopologyId, 0), 10L); - Map taskEndOffsets = Map.of(new TaskId(fooSubtopologyId, 0), 20L); + Map> taskOffsets = Map.of(fooSubtopologyId, Map.of(0, 10L)); + Map> taskEndOffsets = Map.of(fooSubtopologyId, Map.of(0, 20L)); - AssignmentMemberSpec assignmentMemberSpec = createAssignmentMemberSpec( + MemberMetadataAndAssignmentImpl memberMetadata = createMemberMetadataAndAssignment( member, TasksTuple.EMPTY, new MemberTaskOffsets(taskOffsets, taskEndOffsets) ); - assertEquals(taskOffsets, assignmentMemberSpec.taskOffsets()); - assertEquals(taskEndOffsets, assignmentMemberSpec.taskEndOffsets()); + assertEquals(taskOffsets, memberMetadata.taskOffsets()); + assertEquals(taskEndOffsets, memberMetadata.taskEndOffsets()); + } + + @Test + public void testCreateMemberMetadataAndAssignmentSourcesWarmupTasksFromCurrentAssignment() { + String fooSubtopologyId = Uuid.randomUuid().toString(); + + // Warm-up tasks are decided during reconciliation and stored in the member's current + // assignment; they must be sourced from there, not from the target assignment. + StreamsGroupMember member = new StreamsGroupMember.Builder("member-id") + .setInstanceId("instanceId") + .setRackId("rackId") + .setProcessId("processId") + .setClientTags(Map.of()) + .setAssignedTasks(new TasksTupleWithEpochs(Map.of(), Map.of(), Map.of(fooSubtopologyId, Set.of(1, 2, 3)))) + .build(); + + MemberMetadataAndAssignmentImpl memberMetadata = createMemberMetadataAndAssignment( + member, + // The target assignment never carries warm-up tasks; even if it did, it must be ignored. + mkTasksTuple(TaskRole.WARMUP, mkTasks(fooSubtopologyId, 7, 8, 9)), + MemberTaskOffsets.EMPTY + ); + + assertEquals(Map.of(fooSubtopologyId, Set.of(1, 2, 3)), memberMetadata.warmupTasks()); } @Test @@ -167,7 +195,9 @@ public void testEmpty() { @ParameterizedTest - @EnumSource(TaskRole.class) + // Warm-up tasks are not produced by the assignor (only active and standby), so they cannot appear + // in the resulting target assignment. See MemberAssignment. + @EnumSource(value = TaskRole.class, names = {"ACTIVE", "STANDBY"}) public void testAssignmentHasNotChanged(TaskRole taskRole) { TargetAssignmentBuilderTestContext context = new TargetAssignmentBuilderTestContext( "my-group", @@ -221,7 +251,9 @@ public void testAssignmentHasNotChanged(TaskRole taskRole) { @ParameterizedTest - @EnumSource(TaskRole.class) + // Warm-up tasks are not produced by the assignor (only active and standby), so they cannot appear + // in the resulting target assignment. See MemberAssignment. + @EnumSource(value = TaskRole.class, names = {"ACTIVE", "STANDBY"}) public void testAssignmentSwapped(TaskRole taskRole) { TargetAssignmentBuilderTestContext context = new TargetAssignmentBuilderTestContext( "my-group", @@ -288,7 +320,9 @@ public void testAssignmentSwapped(TaskRole taskRole) { @ParameterizedTest - @EnumSource(TaskRole.class) + // Warm-up tasks are not produced by the assignor (only active and standby), so they cannot appear + // in the resulting target assignment. See MemberAssignment. + @EnumSource(value = TaskRole.class, names = {"ACTIVE", "STANDBY"}) public void testPartialAssignmentUpdate(TaskRole taskRole) { TargetAssignmentBuilderTestContext context = new TargetAssignmentBuilderTestContext( "my-group", @@ -404,6 +438,7 @@ public void addGroupMember( memberBuilder.setUserEndpoint(new StreamsGroupMemberMetadataValue.Endpoint().setHost("host").setPort(9090)); memberBuilder.setInstanceId(null); memberBuilder.setRackId(null); + memberBuilder.setAssignedTasks(TasksTupleWithEpochs.EMPTY); members.put(memberId, memberBuilder.build()); targetAssignment.put(memberId, targetTasks); } @@ -424,14 +459,14 @@ public void prepareMemberAssignment( String memberId, TasksTuple assignment ) { - memberAssignments.put(memberId, new MemberAssignment(assignment.activeTasks(), assignment.standbyTasks(), assignment.warmupTasks())); + memberAssignments.put(memberId, new MemberAssignmentImpl(assignment.activeTasks(), assignment.standbyTasks())); } public org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder.TargetAssignmentResult build() { // Prepare expected member specs. - Map memberSpecs = new HashMap<>(); + Map memberMetadataMap = new HashMap<>(); members.forEach((memberId, member) -> - memberSpecs.put(memberId, createAssignmentMemberSpec( + memberMetadataMap.put(memberId, createMemberMetadataAndAssignment( member, targetAssignment.getOrDefault(memberId, TasksTuple.EMPTY), MemberTaskOffsets.EMPTY @@ -444,7 +479,7 @@ public org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder.Target TopologyMetadata topologyMetadata = new TopologyMetadata(metadataImage, subtopologies); // Prepare the expected assignment spec. - GroupSpecImpl groupSpec = new GroupSpecImpl(memberSpecs, new HashMap<>()); + GroupSpecImpl groupSpec = new GroupSpecImpl(memberMetadataMap, new HashMap<>()); // We use `any` here to always return an assignment but use `verify` later on // to ensure that the input was correct. diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImplTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImplTest.java index be2aa696e8f79..1e6e951202418 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImplTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImplTest.java @@ -24,18 +24,21 @@ import java.util.Optional; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; public class GroupSpecImplTest { - private Map members; + private Map members; + private MemberMetadataAndAssignmentImpl member; private GroupSpecImpl groupSpec; @BeforeEach void setUp() { members = new HashMap<>(); - members.put("test-member", new AssignmentMemberSpec( + member = new MemberMetadataAndAssignmentImpl( Optional.of("test-instance"), Optional.of("test-rack"), Map.of(), @@ -45,7 +48,8 @@ void setUp() { Map.of(), Map.of(), Map.of() - )); + ); + members.put("test-member", member); groupSpec = new GroupSpecImpl( members, @@ -54,8 +58,29 @@ void setUp() { } @Test - void testMembers() { - assertEquals(members, groupSpec.members()); + void testMemberIds() { + assertEquals(members.keySet(), groupSpec.memberIds()); + } + + @Test + void testMemberMetadata() { + assertEquals(member, groupSpec.memberMetadata("test-member")); + } + + @Test + void testMemberAssignmentState() { + assertEquals(member, groupSpec.memberAssignmentState("test-member")); + } + + @Test + void testMemberNotFound() { + assertThrows(IllegalArgumentException.class, () -> groupSpec.memberMetadata("unknown")); + assertThrows(IllegalArgumentException.class, () -> groupSpec.memberAssignmentState("unknown")); + } + + @Test + void testConfigs() { + assertTrue(groupSpec.configs().isEmpty()); } } diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/MockAssignorTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/MockAssignorTest.java index da73e27e8c299..79555972e47e1 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/MockAssignorTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/MockAssignorTest.java @@ -16,6 +16,11 @@ */ package org.apache.kafka.coordinator.group.streams.assignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignorException; +import org.apache.kafka.coordinator.group.api.streams.assignor.TopologyDescriber; + import org.junit.jupiter.api.Test; import java.util.HashMap; @@ -60,7 +65,7 @@ public void testZeroMembers() { @Test public void testDoubleAssignment() { - final AssignmentMemberSpec memberSpec1 = new AssignmentMemberSpec( + final MemberMetadataAndAssignmentImpl memberMetadata1 = new MemberMetadataAndAssignmentImpl( Optional.empty(), Optional.empty(), Map.of("test-subtopology", Set.of(0)), @@ -72,7 +77,7 @@ public void testDoubleAssignment() { Map.of() ); - final AssignmentMemberSpec memberSpec2 = new AssignmentMemberSpec( + final MemberMetadataAndAssignmentImpl memberMetadata2 = new MemberMetadataAndAssignmentImpl( Optional.empty(), Optional.empty(), Map.of("test-subtopology", Set.of(0)), @@ -86,7 +91,7 @@ public void testDoubleAssignment() { TaskAssignorException ex = assertThrows(TaskAssignorException.class, () -> assignor.assign( new GroupSpecImpl( - Map.of("member1", memberSpec1, "member2", memberSpec2), + Map.of("member1", memberMetadata1, "member2", memberMetadata2), new HashMap<>() ), new TopologyDescriberImpl(5, List.of("test-subtopology")) @@ -113,7 +118,7 @@ public void testBasicScenario() { @Test public void testSingleMember() { - final AssignmentMemberSpec memberSpec = new AssignmentMemberSpec( + final MemberMetadataAndAssignmentImpl memberMetadata = new MemberMetadataAndAssignmentImpl( Optional.empty(), Optional.empty(), Map.of(), @@ -127,7 +132,7 @@ public void testSingleMember() { final GroupAssignment result = assignor.assign( new GroupSpecImpl( - Map.of("test_member", memberSpec), + Map.of("test_member", memberMetadata), new HashMap<>() ), new TopologyDescriberImpl(4, List.of("test-subtopology")) @@ -145,7 +150,7 @@ public void testSingleMember() { @Test public void testTwoMembersTwoSubtopologies() { - final AssignmentMemberSpec memberSpec1 = new AssignmentMemberSpec( + final MemberMetadataAndAssignmentImpl memberMetadata1 = new MemberMetadataAndAssignmentImpl( Optional.empty(), Optional.empty(), Map.of(), @@ -157,7 +162,7 @@ public void testTwoMembersTwoSubtopologies() { Map.of() ); - final AssignmentMemberSpec memberSpec2 = new AssignmentMemberSpec( + final MemberMetadataAndAssignmentImpl memberMetadata2 = new MemberMetadataAndAssignmentImpl( Optional.empty(), Optional.empty(), Map.of(), @@ -171,7 +176,7 @@ public void testTwoMembersTwoSubtopologies() { final GroupAssignment result = assignor.assign( new GroupSpecImpl( - mkMap(mkEntry("test_member1", memberSpec1), mkEntry("test_member2", memberSpec2)), + mkMap(mkEntry("test_member1", memberMetadata1), mkEntry("test_member2", memberMetadata2)), new HashMap<>() ), new TopologyDescriberImpl(4, List.of("test-subtopology1", "test-subtopology2")) @@ -198,7 +203,7 @@ public void testTwoMembersTwoSubtopologies() { @Test public void testTwoMembersTwoSubtopologiesStickiness() { - final AssignmentMemberSpec memberSpec1 = new AssignmentMemberSpec( + final MemberMetadataAndAssignmentImpl memberMetadata1 = new MemberMetadataAndAssignmentImpl( Optional.empty(), Optional.empty(), mkMap( @@ -213,7 +218,7 @@ public void testTwoMembersTwoSubtopologiesStickiness() { Map.of() ); - final AssignmentMemberSpec memberSpec2 = new AssignmentMemberSpec( + final MemberMetadataAndAssignmentImpl memberMetadata2 = new MemberMetadataAndAssignmentImpl( Optional.empty(), Optional.empty(), mkMap( @@ -229,7 +234,7 @@ public void testTwoMembersTwoSubtopologiesStickiness() { ); final GroupAssignment result = assignor.assign( new GroupSpecImpl( - mkMap(mkEntry("test_member1", memberSpec1), mkEntry("test_member2", memberSpec2)), + mkMap(mkEntry("test_member1", memberMetadata1), mkEntry("test_member2", memberMetadata2)), new HashMap<>() ), new TopologyDescriberImpl(4, List.of("test-subtopology1", "test-subtopology2")) diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignorTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignorTest.java index 161f3e6a5d808..68504a6564c60 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignorTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignorTest.java @@ -17,6 +17,10 @@ package org.apache.kafka.coordinator.group.streams.assignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.TopologyDescriber; + import org.junit.jupiter.api.Test; import org.mockito.internal.util.collections.Sets; @@ -53,13 +57,13 @@ public void testToStringReturnsName() { @Test public void shouldAssignOneActiveTaskToEachProcessWhenTaskCountSameAsProcessCount() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2"); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3"); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2"); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3"); final GroupAssignment result = assignor.assign( new GroupSpecImpl( - mkMap(mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3)), + mkMap(mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)), new HashMap<>() ), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) @@ -78,15 +82,15 @@ public void shouldAssignOneActiveTaskToEachProcessWhenTaskCountSameAsProcessCoun @Test public void shouldAssignTopicGroupIdEvenlyAcrossClientsWithNoStandByTasks() { - final AssignmentMemberSpec memberSpec11 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec12 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec21 = createAssignmentMemberSpec("process2"); - final AssignmentMemberSpec memberSpec22 = createAssignmentMemberSpec("process2"); - final AssignmentMemberSpec memberSpec31 = createAssignmentMemberSpec("process3"); - final AssignmentMemberSpec memberSpec32 = createAssignmentMemberSpec("process3"); - final Map members = mkMap(mkEntry("member1_1", memberSpec11), mkEntry("member1_2", memberSpec12), - mkEntry("member2_1", memberSpec21), mkEntry("member2_2", memberSpec22), - mkEntry("member3_1", memberSpec31), mkEntry("member3_2", memberSpec32)); + final MemberMetadataAndAssignmentImpl memberMetadata11 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata12 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata21 = createMemberMetadata("process2"); + final MemberMetadataAndAssignmentImpl memberMetadata22 = createMemberMetadata("process2"); + final MemberMetadataAndAssignmentImpl memberMetadata31 = createMemberMetadata("process3"); + final MemberMetadataAndAssignmentImpl memberMetadata32 = createMemberMetadata("process3"); + final Map members = mkMap(mkEntry("member1_1", memberMetadata11), mkEntry("member1_2", memberMetadata12), + mkEntry("member2_1", memberMetadata21), mkEntry("member2_2", memberMetadata22), + mkEntry("member3_1", memberMetadata31), mkEntry("member3_2", memberMetadata32)); final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -106,15 +110,15 @@ public void shouldAssignTopicGroupIdEvenlyAcrossClientsWithNoStandByTasks() { @Test public void shouldAssignTopicGroupIdEvenlyAcrossClientsWithStandByTasks() { final Map> tasks = mkMap(mkEntry("test-subtopology1", Sets.newSet(0, 1, 2)), mkEntry("test-subtopology2", Sets.newSet(0, 1, 2))); - final AssignmentMemberSpec memberSpec11 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec12 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec21 = createAssignmentMemberSpec("process2"); - final AssignmentMemberSpec memberSpec22 = createAssignmentMemberSpec("process2"); - final AssignmentMemberSpec memberSpec31 = createAssignmentMemberSpec("process3"); - final AssignmentMemberSpec memberSpec32 = createAssignmentMemberSpec("process3"); - final Map members = mkMap(mkEntry("member1_1", memberSpec11), mkEntry("member1_2", memberSpec12), - mkEntry("member2_1", memberSpec21), mkEntry("member2_2", memberSpec22), - mkEntry("member3_1", memberSpec31), mkEntry("member3_2", memberSpec32)); + final MemberMetadataAndAssignmentImpl memberMetadata11 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata12 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata21 = createMemberMetadata("process2"); + final MemberMetadataAndAssignmentImpl memberMetadata22 = createMemberMetadata("process2"); + final MemberMetadataAndAssignmentImpl memberMetadata31 = createMemberMetadata("process3"); + final MemberMetadataAndAssignmentImpl memberMetadata32 = createMemberMetadata("process3"); + final Map members = mkMap(mkEntry("member1_1", memberMetadata11), mkEntry("member1_2", memberMetadata12), + mkEntry("member2_1", memberMetadata21), mkEntry("member2_2", memberMetadata22), + mkEntry("member3_1", memberMetadata31), mkEntry("member3_2", memberMetadata32)); final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, @@ -135,9 +139,9 @@ public void shouldAssignTopicGroupIdEvenlyAcrossClientsWithStandByTasks() { @Test public void shouldNotMigrateActiveTaskToOtherProcess() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", mkMap(mkEntry("test-subtopology", Set.of(0))), Map.of()); - AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); - Map members = mkMap(mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Set.of(0))), Map.of()); + MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); + Map members = mkMap(mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -154,9 +158,9 @@ public void shouldNotMigrateActiveTaskToOtherProcess() { testMember1.activeTasks().get("test-subtopology").size() + testMember2.activeTasks().get("test-subtopology").size()); // flip the previous active tasks assignment around. - memberSpec2 = createAssignmentMemberSpec("process2", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process1", mkMap(mkEntry("test-subtopology", Set.of(2))), Map.of()); - members = mkMap(mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3)); + memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Set.of(2))), Map.of()); + members = mkMap(mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -175,11 +179,11 @@ public void shouldNotMigrateActiveTaskToOtherProcess() { @Test public void shouldMigrateActiveTasksToNewProcessWithoutChangingAllAssignments() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 2))), Map.of()); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3"); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 2))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3"); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -203,11 +207,11 @@ public void shouldMigrateActiveTasksToNewProcessWithoutChangingAllAssignments() @Test public void shouldAssignBasedOnCapacity() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec21 = createAssignmentMemberSpec("process2"); - final AssignmentMemberSpec memberSpec22 = createAssignmentMemberSpec("process2"); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2_1", memberSpec21), mkEntry("member2_2", memberSpec22)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata21 = createMemberMetadata("process2"); + final MemberMetadataAndAssignmentImpl memberMetadata22 = createMemberMetadata("process2"); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2_1", memberMetadata21), mkEntry("member2_2", memberMetadata22)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -230,10 +234,10 @@ public void shouldAssignTasksEvenlyWithUnequalTopicGroupSizes() { final Map> activeTasks = mkMap( mkEntry("test-subtopology1", Sets.newSet(0, 1, 2, 3, 4, 5)), mkEntry("test-subtopology2", Sets.newSet(0))); - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", activeTasks, Map.of()); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2"); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", activeTasks, Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2"); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -257,14 +261,14 @@ public void shouldAssignTasksEvenlyWithUnequalTopicGroupSizes() { @Test public void shouldKeepActiveTaskStickinessWhenMoreClientThanActiveTasks() { - AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", mkMap(mkEntry("test-subtopology", Set.of(0))), Map.of()); - AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", mkMap(mkEntry("test-subtopology", Set.of(2))), Map.of()); - AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); - AssignmentMemberSpec memberSpec4 = createAssignmentMemberSpec("process4"); - AssignmentMemberSpec memberSpec5 = createAssignmentMemberSpec("process5"); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), - mkEntry("member3", memberSpec3), mkEntry("member4", memberSpec4), mkEntry("member5", memberSpec5)); + MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Set.of(0))), Map.of()); + MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Set.of(2))), Map.of()); + MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); + MemberMetadataAndAssignmentImpl memberMetadata4 = createMemberMetadata("process4"); + MemberMetadataAndAssignmentImpl memberMetadata5 = createMemberMetadata("process5"); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), + mkEntry("member3", memberMetadata3), mkEntry("member4", memberMetadata4), mkEntry("member5", memberMetadata5)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -291,14 +295,14 @@ public void shouldKeepActiveTaskStickinessWhenMoreClientThanActiveTasks() { assertNull(testMember5.activeTasks().get("test-subtopology")); // change up the assignment and make sure it is still sticky - memberSpec1 = createAssignmentMemberSpec("process1"); - memberSpec2 = createAssignmentMemberSpec("process2", mkMap(mkEntry("test-subtopology", Set.of(0))), Map.of()); - memberSpec3 = createAssignmentMemberSpec("process3"); - memberSpec4 = createAssignmentMemberSpec("process4", mkMap(mkEntry("test-subtopology", Set.of(2))), Map.of()); - memberSpec5 = createAssignmentMemberSpec("process5", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); + memberMetadata1 = createMemberMetadata("process1"); + memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Set.of(0))), Map.of()); + memberMetadata3 = createMemberMetadata("process3"); + memberMetadata4 = createMemberMetadata("process4", mkMap(mkEntry("test-subtopology", Set.of(2))), Map.of()); + memberMetadata5 = createMemberMetadata("process5", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), - mkEntry("member3", memberSpec3), mkEntry("member4", memberSpec4), mkEntry("member5", memberSpec5)); + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), + mkEntry("member3", memberMetadata3), mkEntry("member4", memberMetadata4), mkEntry("member5", memberMetadata5)); result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -327,11 +331,11 @@ public void shouldKeepActiveTaskStickinessWhenMoreClientThanActiveTasks() { @Test public void shouldAssignTasksToClientWithPreviousStandbyTasks() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", Map.of(), mkMap(mkEntry("test-subtopology", Set.of(2)))); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", Map.of(), mkMap(mkEntry("test-subtopology", Set.of(1)))); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3", Map.of(), mkMap(mkEntry("test-subtopology", Set.of(0)))); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", Map.of(), mkMap(mkEntry("test-subtopology", Set.of(2)))); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", Map.of(), mkMap(mkEntry("test-subtopology", Set.of(1)))); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3", Map.of(), mkMap(mkEntry("test-subtopology", Set.of(0)))); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -354,10 +358,10 @@ public void shouldAssignTasksToClientWithPreviousStandbyTasks() { @Test public void shouldNotAssignStandbyTasksToClientWithPreviousStandbyTasksAndCurrentActiveTasks() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", Map.of(), mkMap(mkEntry("test-subtopology", Set.of(0)))); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", Map.of(), mkMap(mkEntry("test-subtopology", Set.of(1)))); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", Map.of(), mkMap(mkEntry("test-subtopology", Set.of(0)))); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", Map.of(), mkMap(mkEntry("test-subtopology", Set.of(1)))); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, "1"))), @@ -380,17 +384,17 @@ public void shouldNotAssignStandbyTasksToClientWithPreviousStandbyTasksAndCurren @Test public void shouldAssignBasedOnCapacityWhenMultipleClientHaveStandbyTasks() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Set.of(0))), mkMap(mkEntry("test-subtopology", Set.of(1)))); - final AssignmentMemberSpec memberSpec21 = createAssignmentMemberSpec("process2", + final MemberMetadataAndAssignmentImpl memberMetadata21 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Set.of(2))), mkMap(mkEntry("test-subtopology", Set.of(1)))); - final AssignmentMemberSpec memberSpec22 = createAssignmentMemberSpec("process2", + final MemberMetadataAndAssignmentImpl memberMetadata22 = createMemberMetadata("process2", Map.of(), Map.of()); - Map members = mkMap( - mkEntry("member1", memberSpec1), - mkEntry("member2_1", memberSpec21), mkEntry("member2_2", memberSpec22)); + Map members = mkMap( + mkEntry("member1", memberMetadata1), + mkEntry("member2_1", memberMetadata21), mkEntry("member2_2", memberMetadata22)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -414,13 +418,13 @@ public void shouldAssignBasedOnCapacityWhenMultipleClientHaveStandbyTasks() { @Test public void shouldAssignStandbyTasksToDifferentClientThanCorrespondingActiveTaskIsAssignedTo() { final Map> tasks = mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1, 2, 3))); - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", mkMap(mkEntry("test-subtopology", Set.of(0))), Map.of()); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3", mkMap(mkEntry("test-subtopology", Set.of(2))), Map.of()); - final AssignmentMemberSpec memberSpec4 = createAssignmentMemberSpec("process4", mkMap(mkEntry("test-subtopology", Set.of(3))), Map.of()); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), - mkEntry("member3", memberSpec3), mkEntry("member4", memberSpec4)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Set.of(0))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3", mkMap(mkEntry("test-subtopology", Set.of(2))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata4 = createMemberMetadata("process4", mkMap(mkEntry("test-subtopology", Set.of(3))), Map.of()); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), + mkEntry("member3", memberMetadata3), mkEntry("member4", memberMetadata4)); final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, @@ -451,12 +455,12 @@ public void shouldAssignStandbyTasksToDifferentClientThanCorrespondingActiveTask @Test public void shouldAssignMultipleReplicasOfStandbyTask() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", mkMap(mkEntry("test-subtopology", Set.of(0))), Map.of()); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3", mkMap(mkEntry("test-subtopology", Set.of(2))), Map.of()); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), - mkEntry("member3", memberSpec3)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Set.of(0))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3", mkMap(mkEntry("test-subtopology", Set.of(2))), Map.of()); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), + mkEntry("member3", memberMetadata3)); final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, @@ -471,9 +475,9 @@ public void shouldAssignMultipleReplicasOfStandbyTask() { @Test public void shouldNotAssignStandbyTaskReplicasWhenNoClientAvailableWithoutHavingTheTaskAssigned() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1"); - Map members = mkMap( - mkEntry("member1", memberSpec1)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1"); + Map members = mkMap( + mkEntry("member1", memberMetadata1)); final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, @@ -486,12 +490,12 @@ public void shouldNotAssignStandbyTaskReplicasWhenNoClientAvailableWithoutHaving @Test public void shouldAssignActiveAndStandbyTasks() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2"); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3"); - Map members = mkMap( - mkEntry("member1", memberSpec1), - mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2"); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3"); + Map members = mkMap( + mkEntry("member1", memberMetadata1), + mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, @@ -505,14 +509,14 @@ public void shouldAssignActiveAndStandbyTasks() { @Test public void shouldAssignAtLeastOneTaskToEachClientIfPossible() { - final AssignmentMemberSpec memberSpec11 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec12 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec13 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2"); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3"); - Map members = mkMap( - mkEntry("member1_1", memberSpec11), mkEntry("member1_2", memberSpec12), mkEntry("member1_3", memberSpec13), - mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3)); + final MemberMetadataAndAssignmentImpl memberMetadata11 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata12 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata13 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2"); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3"); + Map members = mkMap( + mkEntry("member1_1", memberMetadata11), mkEntry("member1_2", memberMetadata12), mkEntry("member1_3", memberMetadata13), + mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -526,15 +530,15 @@ public void shouldAssignAtLeastOneTaskToEachClientIfPossible() { @Test public void shouldAssignEachActiveTaskToOneClientWhenMoreClientsThanTasks() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2"); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3"); - final AssignmentMemberSpec memberSpec4 = createAssignmentMemberSpec("process4"); - final AssignmentMemberSpec memberSpec5 = createAssignmentMemberSpec("process5"); - final AssignmentMemberSpec memberSpec6 = createAssignmentMemberSpec("process6"); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3), - mkEntry("member4", memberSpec4), mkEntry("member5", memberSpec5), mkEntry("member6", memberSpec6)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2"); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3"); + final MemberMetadataAndAssignmentImpl memberMetadata4 = createMemberMetadata("process4"); + final MemberMetadataAndAssignmentImpl memberMetadata5 = createMemberMetadata("process5"); + final MemberMetadataAndAssignmentImpl memberMetadata6 = createMemberMetadata("process6"); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3), + mkEntry("member4", memberMetadata4), mkEntry("member5", memberMetadata5), mkEntry("member6", memberMetadata6)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -547,15 +551,15 @@ public void shouldAssignEachActiveTaskToOneClientWhenMoreClientsThanTasks() { @Test public void shouldBalanceActiveAndStandbyTasksAcrossAvailableClients() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2"); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3"); - final AssignmentMemberSpec memberSpec4 = createAssignmentMemberSpec("process4"); - final AssignmentMemberSpec memberSpec5 = createAssignmentMemberSpec("process5"); - final AssignmentMemberSpec memberSpec6 = createAssignmentMemberSpec("process6"); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3), - mkEntry("member4", memberSpec4), mkEntry("member5", memberSpec5), mkEntry("member6", memberSpec6)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2"); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3"); + final MemberMetadataAndAssignmentImpl memberMetadata4 = createMemberMetadata("process4"); + final MemberMetadataAndAssignmentImpl memberMetadata5 = createMemberMetadata("process5"); + final MemberMetadataAndAssignmentImpl memberMetadata6 = createMemberMetadata("process6"); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3), + mkEntry("member4", memberMetadata4), mkEntry("member5", memberMetadata5), mkEntry("member6", memberMetadata6)); final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, @@ -570,11 +574,11 @@ public void shouldBalanceActiveAndStandbyTasksAcrossAvailableClients() { @Test public void shouldAssignMoreTasksToClientWithMoreCapacity() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec21 = createAssignmentMemberSpec("process2"); - final AssignmentMemberSpec memberSpec22 = createAssignmentMemberSpec("process2"); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2_1", memberSpec21), mkEntry("member2_2", memberSpec22)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata21 = createMemberMetadata("process2"); + final MemberMetadataAndAssignmentImpl memberMetadata22 = createMemberMetadata("process2"); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2_1", memberMetadata21), mkEntry("member2_2", memberMetadata22)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -587,12 +591,12 @@ public void shouldAssignMoreTasksToClientWithMoreCapacity() { @Test public void shouldReBalanceTasksAcrossAllClientsWhenCapacityAndTaskCountTheSame() { - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1, 2, 3))), Map.of()); - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2"); - final AssignmentMemberSpec memberSpec4 = createAssignmentMemberSpec("process4"); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3), mkEntry("member4", memberSpec4)); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1, 2, 3))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2"); + final MemberMetadataAndAssignmentImpl memberMetadata4 = createMemberMetadata("process4"); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3), mkEntry("member4", memberMetadata4)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -607,11 +611,11 @@ public void shouldReBalanceTasksAcrossAllClientsWhenCapacityAndTaskCountTheSame( @Test public void shouldReBalanceTasksAcrossClientsWhenCapacityLessThanTaskCount() { - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1, 2, 3))), Map.of()); - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2"); - Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3)); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1, 2, 3))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2"); + Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -625,11 +629,11 @@ public void shouldReBalanceTasksAcrossClientsWhenCapacityLessThanTaskCount() { @Test public void shouldRebalanceTasksToClientsBasedOnCapacity() { - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 3, 2))), Map.of()); - final AssignmentMemberSpec memberSpec31 = createAssignmentMemberSpec("process3"); - final AssignmentMemberSpec memberSpec32 = createAssignmentMemberSpec("process3"); - Map members = mkMap( - mkEntry("member2", memberSpec2), mkEntry("member3_1", memberSpec31), mkEntry("member3_2", memberSpec32)); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 3, 2))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata31 = createMemberMetadata("process3"); + final MemberMetadataAndAssignmentImpl memberMetadata32 = createMemberMetadata("process3"); + Map members = mkMap( + mkEntry("member2", memberMetadata2), mkEntry("member3_1", memberMetadata31), mkEntry("member3_2", memberMetadata32)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -644,11 +648,11 @@ public void shouldRebalanceTasksToClientsBasedOnCapacity() { public void shouldMoveMinimalNumberOfTasksWhenPreviouslyAboveCapacityAndNewClientAdded() { final Set p1PrevTasks = Sets.newSet(0, 2); final Set p2PrevTasks = Sets.newSet(1, 3); - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", mkMap(mkEntry("test-subtopology", p1PrevTasks)), Map.of()); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", mkMap(mkEntry("test-subtopology", p2PrevTasks)), Map.of()); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3"); - final Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", p1PrevTasks)), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", p2PrevTasks)), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3"); + final Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -666,10 +670,10 @@ public void shouldMoveMinimalNumberOfTasksWhenPreviouslyAboveCapacityAndNewClien @Test public void shouldNotMoveAnyTasksWhenNewTasksAdded() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1))), Map.of()); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", mkMap(mkEntry("test-subtopology", Sets.newSet(2, 3))), Map.of()); - final Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Sets.newSet(2, 3))), Map.of()); + final Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -686,11 +690,11 @@ public void shouldNotMoveAnyTasksWhenNewTasksAdded() { @Test public void shouldAssignNewTasksToNewClientWhenPreviousTasksAssignedToOldClients() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(2, 1))), Map.of()); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 3))), Map.of()); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3"); - final Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(2, 1))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 3))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3"); + final Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -710,20 +714,20 @@ public void shouldAssignNewTasksToNewClientWhenPreviousTasksAssignedToOldClients @Test public void shouldAssignTasksNotPreviouslyActiveToNewClient() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology0", Sets.newSet(1)), mkEntry("test-subtopology1", Sets.newSet(2, 3))), mkMap(mkEntry("test-subtopology0", Sets.newSet(0)), mkEntry("test-subtopology1", Sets.newSet(1)), mkEntry("test-subtopology2", Sets.newSet(0, 1, 3)))); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology0", Sets.newSet(0)), mkEntry("test-subtopology1", Sets.newSet(1)), mkEntry("test-subtopology2", Sets.newSet(2))), mkMap(mkEntry("test-subtopology0", Sets.newSet(1, 2, 3)), mkEntry("test-subtopology1", Sets.newSet(0, 2, 3)), mkEntry("test-subtopology2", Sets.newSet(0, 1, 3)))); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3", + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3", mkMap(mkEntry("test-subtopology2", Sets.newSet(0, 1, 3))), mkMap(mkEntry("test-subtopology0", Sets.newSet(2)), mkEntry("test-subtopology1", Sets.newSet(2)))); - final AssignmentMemberSpec newMemberSpec = createAssignmentMemberSpec("process4", + final MemberMetadataAndAssignmentImpl newMemberSpec = createMemberMetadata("process4", Map.of(), mkMap(mkEntry("test-subtopology0", Sets.newSet(0, 1, 2, 3)), mkEntry("test-subtopology1", Sets.newSet(0, 1, 2, 3)), mkEntry("test-subtopology2", Sets.newSet(0, 1, 2, 3)))); - final Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3), mkEntry("newMember", newMemberSpec)); + final Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3), mkEntry("newMember", newMemberSpec)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -742,20 +746,20 @@ public void shouldAssignTasksNotPreviouslyActiveToNewClient() { @Test public void shouldAssignTasksNotPreviouslyActiveToMultipleNewClients() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology0", Sets.newSet(1)), mkEntry("test-subtopology1", Sets.newSet(2, 3))), mkMap(mkEntry("test-subtopology0", Sets.newSet(0)), mkEntry("test-subtopology1", Sets.newSet(1)), mkEntry("test-subtopology2", Sets.newSet(0, 1, 3)))); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology0", Sets.newSet(0)), mkEntry("test-subtopology1", Sets.newSet(1)), mkEntry("test-subtopology2", Sets.newSet(2))), mkMap(mkEntry("test-subtopology0", Sets.newSet(1, 2, 3)), mkEntry("test-subtopology1", Sets.newSet(0, 2, 3)), mkEntry("test-subtopology2", Sets.newSet(0, 1, 3)))); - final AssignmentMemberSpec bounce1 = createAssignmentMemberSpec("bounce1", + final MemberMetadataAndAssignmentImpl bounce1 = createMemberMetadata("bounce1", Map.of(), mkMap(mkEntry("test-subtopology2", Sets.newSet(0, 1, 3)))); - final AssignmentMemberSpec bounce2 = createAssignmentMemberSpec("bounce2", + final MemberMetadataAndAssignmentImpl bounce2 = createMemberMetadata("bounce2", Map.of(), mkMap(mkEntry("test-subtopology0", Sets.newSet(2, 3)), mkEntry("test-subtopology1", Sets.newSet(0)))); - final Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("bounce_member1", bounce1), mkEntry("bounce_member2", bounce2)); + final Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("bounce_member1", bounce1), mkEntry("bounce_member2", bounce2)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -774,10 +778,10 @@ public void shouldAssignTasksNotPreviouslyActiveToMultipleNewClients() { @Test public void shouldAssignTasksToNewClient() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(1, 2))), Map.of()); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2"); - final Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(1, 2))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2"); + final Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -789,11 +793,11 @@ public void shouldAssignTasksToNewClient() { @Test public void shouldAssignTasksToNewClientWithoutFlippingAssignmentBetweenExistingClients() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1, 2))), Map.of()); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", mkMap(mkEntry("test-subtopology", Sets.newSet(3, 4, 5))), Map.of()); - final AssignmentMemberSpec newMemberSpec = createAssignmentMemberSpec("process3"); - final Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("newMember", newMemberSpec)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1, 2))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Sets.newSet(3, 4, 5))), Map.of()); + final MemberMetadataAndAssignmentImpl newMemberSpec = createMemberMetadata("process3"); + final Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("newMember", newMemberSpec)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -815,11 +819,11 @@ public void shouldAssignTasksToNewClientWithoutFlippingAssignmentBetweenExisting @Test public void shouldAssignTasksToNewClientWithoutFlippingAssignmentBetweenExistingAndBouncedClients() { - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1, 2, 6))), Map.of()); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", Map.of(), mkMap(mkEntry("test-subtopology", Sets.newSet(3, 4, 5)))); - final AssignmentMemberSpec newMemberSpec = createAssignmentMemberSpec("newProcess"); - final Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("newMember", newMemberSpec)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1, 2, 6))), Map.of()); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", Map.of(), mkMap(mkEntry("test-subtopology", Sets.newSet(3, 4, 5)))); + final MemberMetadataAndAssignmentImpl newMemberSpec = createMemberMetadata("newProcess"); + final Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("newMember", newMemberSpec)); GroupAssignment result = assignor.assign( new GroupSpecImpl(members, new HashMap<>()), @@ -845,9 +849,9 @@ public void shouldHandleLargeNumberOfTasksWithStandbyAssignment() { final int numClients = 5; final int numStandbyReplicas = 2; - Map members = new HashMap<>(); + Map members = new HashMap<>(); for (int i = 0; i < numClients; i++) { - members.put("member" + i, createAssignmentMemberSpec("process" + i)); + members.put("member" + i, createMemberMetadata("process" + i)); } GroupAssignment result = assignor.assign( @@ -916,9 +920,9 @@ public void shouldHandleOddNumberOfClientsWithStandbyTasks() { final int numClients = 7; final int numStandbyReplicas = 1; - Map members = new HashMap<>(); + Map members = new HashMap<>(); for (int i = 0; i < numClients; i++) { - members.put("member" + i, createAssignmentMemberSpec("process" + i)); + members.put("member" + i, createMemberMetadata("process" + i)); } GroupAssignment result = assignor.assign( @@ -967,9 +971,9 @@ public void shouldHandleHighStandbyReplicaCount() { final int numClients = 3; final int numStandbyReplicas = 5; - Map members = new HashMap<>(); + Map members = new HashMap<>(); for (int i = 0; i < numClients; i++) { - members.put("member" + i, createAssignmentMemberSpec("process" + i)); + members.put("member" + i, createMemberMetadata("process" + i)); } GroupAssignment result = assignor.assign( @@ -1011,9 +1015,9 @@ public void shouldHandleLargeNumberOfSubtopologiesWithStandbyTasks() { subtopologies.add("subtopology-" + i); } - Map members = new HashMap<>(); + Map members = new HashMap<>(); for (int i = 0; i < numClients; i++) { - members.put("member" + i, createAssignmentMemberSpec("process" + i)); + members.put("member" + i, createMemberMetadata("process" + i)); } GroupAssignment result = assignor.assign( @@ -1044,8 +1048,8 @@ public void shouldHandleEdgeCaseWithSingleClientAndMultipleStandbyReplicas() { final int numTasks = 10; final int numStandbyReplicas = 3; - Map members = mkMap( - mkEntry("member1", createAssignmentMemberSpec("process1")) + Map members = mkMap( + mkEntry("member1", createMemberMetadata("process1")) ); GroupAssignment result = assignor.assign( @@ -1068,9 +1072,9 @@ public void shouldHandleEdgeCaseWithMoreStandbyReplicasThanAvailableClients() { final int numClients = 2; final int numStandbyReplicas = 5; // More than available clients - Map members = new HashMap<>(); + Map members = new HashMap<>(); for (int i = 0; i < numClients; i++) { - members.put("member" + i, createAssignmentMemberSpec("process" + i)); + members.put("member" + i, createMemberMetadata("process" + i)); } GroupAssignment result = assignor.assign( @@ -1101,19 +1105,19 @@ public void shouldHandleEdgeCaseWithMoreStandbyReplicasThanAvailableClients() { public void shouldReassignTasksWhenNewNodeJoinsWithExistingActiveAndStandbyAssignments() { // Initial setup: Node 1 has active tasks 0,1 and standby tasks 2,3 // Node 2 has active tasks 2,3 and standby tasks 0,1 - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1))), mkMap(mkEntry("test-subtopology", Sets.newSet(2, 3)))); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Sets.newSet(2, 3))), mkMap(mkEntry("test-subtopology", Sets.newSet(0, 1)))); // Node 3 joins as new client - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3"); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3"); - final Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2), mkEntry("member3", memberSpec3)); + final Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, "1"))), @@ -1148,12 +1152,12 @@ public void shouldReassignTasksWhenNewNodeJoinsWithExistingActiveAndStandbyAssig @Test public void shouldRangeAssignTasksWhenScalingUp() { // Two clients, the second one is new - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", Map.of("test-subtopology1", Set.of(0, 1), "test-subtopology2", Set.of(0, 1)), Map.of()); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2"); - final Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2)); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2"); + final Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2)); // Two subtopologies with 2 tasks each (4 tasks total) with standby replicas enabled final GroupAssignment result = assignor.assign( @@ -1196,10 +1200,10 @@ public void shouldRangeAssignTasksWhenScalingUp() { @Test public void shouldRangeAssignTasksWhenStartingEmpty() { // Two clients starting empty (no previous tasks) - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1"); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2"); - final Map members = mkMap( - mkEntry("member1", memberSpec1), mkEntry("member2", memberSpec2)); + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1"); + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2"); + final Map members = mkMap( + mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2)); // Two subtopologies with 2 tasks each (4 tasks total) with standby replicas enabled final GroupAssignment result = assignor.assign( @@ -1257,18 +1261,18 @@ public void shouldAssignStandbyTaskToPreviousOwnerBasedOnBelowQuotaCondition() { // Process1: active=[0], standby=[1] (previously had both active and standby tasks) // Process2: active=[1] (had the active task that process1 had as standby) // Process3: no previous tasks - final AssignmentMemberSpec memberSpec1 = createAssignmentMemberSpec("process1", + final MemberMetadataAndAssignmentImpl memberMetadata1 = createMemberMetadata("process1", mkMap(mkEntry("test-subtopology", Sets.newSet(0))), mkMap(mkEntry("test-subtopology", Sets.newSet(1)))); - final AssignmentMemberSpec memberSpec2 = createAssignmentMemberSpec("process2", + final MemberMetadataAndAssignmentImpl memberMetadata2 = createMemberMetadata("process2", mkMap(mkEntry("test-subtopology", Sets.newSet(1))), Map.of()); - final AssignmentMemberSpec memberSpec3 = createAssignmentMemberSpec("process3"); + final MemberMetadataAndAssignmentImpl memberMetadata3 = createMemberMetadata("process3"); - final Map members = mkMap( - mkEntry("member1", memberSpec1), - mkEntry("member2", memberSpec2), - mkEntry("member3", memberSpec3)); + final Map members = mkMap( + mkEntry("member1", memberMetadata1), + mkEntry("member2", memberMetadata2), + mkEntry("member3", memberMetadata3)); // We have 2 active tasks + 1 standby replica = 4 total tasks // Quota per process = 4 tasks / 3 processes = 1.33 -> 2 tasks per process @@ -1429,8 +1433,8 @@ private Map> mergeAllStandbyTasks(GroupAssignment result) { return mergeAllStandbyTasks(result, result.members().keySet().toArray(memberIds)); } - private AssignmentMemberSpec createAssignmentMemberSpec(final String processId) { - return new AssignmentMemberSpec( + private MemberMetadataAndAssignmentImpl createMemberMetadata(final String processId) { + return new MemberMetadataAndAssignmentImpl( Optional.empty(), Optional.empty(), Map.of(), @@ -1442,9 +1446,9 @@ private AssignmentMemberSpec createAssignmentMemberSpec(final String processId) Map.of()); } - private AssignmentMemberSpec createAssignmentMemberSpec(final String processId, final Map> prevActiveTasks, + private MemberMetadataAndAssignmentImpl createMemberMetadata(final String processId, final Map> prevActiveTasks, final Map> prevStandbyTasks) { - return new AssignmentMemberSpec( + return new MemberMetadataAndAssignmentImpl( Optional.empty(), Optional.empty(), prevActiveTasks, diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/StreamsAssignorBenchmarkUtils.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/StreamsAssignorBenchmarkUtils.java index fa3324d012e64..91a48662f2a16 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/StreamsAssignorBenchmarkUtils.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/StreamsAssignorBenchmarkUtils.java @@ -16,10 +16,10 @@ */ package org.apache.kafka.jmh.assignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupSpec; import org.apache.kafka.coordinator.group.streams.StreamsGroupMember; -import org.apache.kafka.coordinator.group.streams.assignor.AssignmentMemberSpec; -import org.apache.kafka.coordinator.group.streams.assignor.GroupSpec; import org.apache.kafka.coordinator.group.streams.assignor.GroupSpecImpl; +import org.apache.kafka.coordinator.group.streams.assignor.MemberMetadataAndAssignmentImpl; import org.apache.kafka.coordinator.group.streams.topics.ConfiguredInternalTopic; import org.apache.kafka.coordinator.group.streams.topics.ConfiguredSubtopology; @@ -46,14 +46,14 @@ public static GroupSpec createGroupSpec( Map members, Map assignmentConfigs ) { - Map memberSpecs = new HashMap<>(); + Map memberSpecs = new HashMap<>(); // Prepare the member spec for all members. for (Map.Entry memberEntry : members.entrySet()) { String memberId = memberEntry.getKey(); StreamsGroupMember member = memberEntry.getValue(); - memberSpecs.put(memberId, new AssignmentMemberSpec( + memberSpecs.put(memberId, new MemberMetadataAndAssignmentImpl( member.instanceId(), member.rackId(), Map.of(), diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/StreamsStickyAssignorBenchmark.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/StreamsStickyAssignorBenchmark.java index 22863fef5fa2e..b51a776ff3cfc 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/StreamsStickyAssignorBenchmark.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/StreamsStickyAssignorBenchmark.java @@ -17,16 +17,17 @@ package org.apache.kafka.jmh.assignor; import org.apache.kafka.coordinator.common.runtime.CoordinatorMetadataImage; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.GroupSpec; +import org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignment; +import org.apache.kafka.coordinator.group.api.streams.assignor.TaskAssignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.TopologyDescriber; import org.apache.kafka.coordinator.group.streams.StreamsGroupMember; import org.apache.kafka.coordinator.group.streams.TopologyMetadata; -import org.apache.kafka.coordinator.group.streams.assignor.AssignmentMemberSpec; -import org.apache.kafka.coordinator.group.streams.assignor.GroupAssignment; -import org.apache.kafka.coordinator.group.streams.assignor.GroupSpec; import org.apache.kafka.coordinator.group.streams.assignor.GroupSpecImpl; -import org.apache.kafka.coordinator.group.streams.assignor.MemberAssignment; +import org.apache.kafka.coordinator.group.streams.assignor.MemberAssignmentImpl; +import org.apache.kafka.coordinator.group.streams.assignor.MemberMetadataAndAssignmentImpl; import org.apache.kafka.coordinator.group.streams.assignor.StickyTaskAssignor; -import org.apache.kafka.coordinator.group.streams.assignor.TaskAssignor; -import org.apache.kafka.coordinator.group.streams.assignor.TopologyDescriber; import org.apache.kafka.coordinator.group.streams.topics.ConfiguredSubtopology; import org.openjdk.jmh.annotations.Benchmark; @@ -132,28 +133,29 @@ private void simulateIncrementalRebalance() { GroupAssignment initialAssignment = new StickyTaskAssignor().assign(groupSpec, topologyDescriber); Map members = initialAssignment.members(); - Map updatedMemberSpec = new HashMap<>(); + Map updatedMemberSpec = new HashMap<>(); - for (Map.Entry member : groupSpec.members().entrySet()) { + for (String memberId : groupSpec.memberIds()) { MemberAssignment memberAssignment = members.getOrDefault( - member.getKey(), - new MemberAssignment(Map.of(), Map.of(), Map.of()) + memberId, + new MemberAssignmentImpl(Map.of(), Map.of()) ); - updatedMemberSpec.put(member.getKey(), new AssignmentMemberSpec( + updatedMemberSpec.put(memberId, new MemberMetadataAndAssignmentImpl( Optional.empty(), Optional.empty(), memberAssignment.activeTasks(), memberAssignment.standbyTasks(), - memberAssignment.warmupTasks(), - member.getValue().processId(), + // Warm-up tasks are not assigned by the assignor; they are decided during reconciliation. + Map.of(), + groupSpec.memberMetadata(memberId).processId(), Map.of(), Map.of(), Map.of() )); } - updatedMemberSpec.put("newMember", new AssignmentMemberSpec( + updatedMemberSpec.put("newMember", new MemberMetadataAndAssignmentImpl( Optional.empty(), Optional.empty(), Map.of(),