diff --git a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/AssignmentConfigs.java b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/AssignmentConfigs.java new file mode 100644 index 0000000000000..3b416ff93c91f --- /dev/null +++ b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/AssignmentConfigs.java @@ -0,0 +1,43 @@ +/* + * 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.List; + +/** + * The assignment configurations that the group coordinator passes to the task assignor. + * + *

This interface is not intended to be implemented by task assignors: new configurations may be added to it. + */ +@InterfaceAudience.Public +@InterfaceStability.Evolving +public interface AssignmentConfigs { + + /** + * @return The number of standby replicas for each task. + */ + int numStandbyReplicas(); + + /** + * @return The client tags used to distribute standby tasks across racks. The list is unmodifiable. + */ + List rackAwareAssignmentTags(); + +} 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 index 35aeb291bdf3c..03485c14bfb31 100644 --- 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 @@ -20,7 +20,6 @@ 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. @@ -51,8 +50,8 @@ public interface GroupSpec { MemberAssignmentState memberAssignmentState(String memberId); /** - * @return Any configurations passed to the assignor. The map is unmodifiable. + * @return The assignment configurations passed to the assignor. */ - Map configs(); + AssignmentConfigs configs(); } 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 124c392ba42de..32036f04502fe 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,10 +19,12 @@ 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.api.streams.assignor.AssignmentConfigs; 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.AssignmentConfigsImpl; import org.apache.kafka.coordinator.group.streams.assignor.GroupSpecImpl; import org.apache.kafka.coordinator.group.streams.assignor.MemberMetadataAndStateImpl; import org.apache.kafka.coordinator.group.streams.topics.ConfiguredTopology; @@ -70,7 +72,7 @@ public class TargetAssignmentBuilder { /** * The assignment configs. */ - private final Map assignmentConfigs; + private final AssignmentConfigs assignmentConfigs; /** * The members in the group. @@ -114,7 +116,7 @@ public TargetAssignmentBuilder( this.groupId = Objects.requireNonNull(groupId); this.groupEpoch = groupEpoch; this.assignor = Objects.requireNonNull(assignor); - this.assignmentConfigs = Objects.requireNonNull(assignmentConfigs); + this.assignmentConfigs = AssignmentConfigsImpl.fromMap(Objects.requireNonNull(assignmentConfigs)); } static MemberMetadataAndStateImpl createMemberMetadataAndState( diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/AssignmentConfigsImpl.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/AssignmentConfigsImpl.java new file mode 100644 index 0000000000000..437244ddbdc98 --- /dev/null +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/AssignmentConfigsImpl.java @@ -0,0 +1,67 @@ +/* + * 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.AssignmentConfigs; + +import java.util.List; +import java.util.Map; +import java.util.Objects; + +/** + * The assignment configurations for a streams group. + * + * @param numStandbyReplicas The number of standby replicas for each task. + * @param rackAwareAssignmentTags The client tags used to distribute standby tasks across racks. + */ +public record AssignmentConfigsImpl( + int numStandbyReplicas, + List rackAwareAssignmentTags +) implements AssignmentConfigs { + + private static final String NUM_STANDBY_REPLICAS_CONFIG = "num.standby.replicas"; + private static final String RACK_AWARE_ASSIGNMENT_TAGS_CONFIG = "rack.aware.assignment.tags"; + + /** + * The configs used for a group that has none of them set. + */ + public static final AssignmentConfigsImpl DEFAULT = new AssignmentConfigsImpl(0, List.of()); + + public AssignmentConfigsImpl { + // The list is exposed to a custom assignor through the public AssignmentConfigs interface. + rackAwareAssignmentTags = List.copyOf(Objects.requireNonNull(rackAwareAssignmentTags)); + } + + /** + * Converts the raw assignment configs computed for the group into the typed configs passed to the assignor. + */ + public static AssignmentConfigsImpl fromMap(Map configs) { + // The map is empty when it was replayed from a group metadata record written before the last assignment + // configs were persisted. + if (configs.isEmpty()) { + return DEFAULT; + } + // The rack-aware assignment tags are only set when any are configured. + String rackAwareAssignmentTags = configs.get(RACK_AWARE_ASSIGNMENT_TAGS_CONFIG); + return new AssignmentConfigsImpl( + Integer.parseInt(configs.get(NUM_STANDBY_REPLICAS_CONFIG)), + rackAwareAssignmentTags == null + ? List.of() + : List.of(rackAwareAssignmentTags.trim().split("\\s*,\\s*", -1)) + ); + } +} 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 19b4d08556bfd..ec408f65a9f48 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,6 +16,7 @@ */ package org.apache.kafka.coordinator.group.streams.assignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.AssignmentConfigs; 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; @@ -30,17 +31,17 @@ * * @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. + * @param configs The assignment configurations passed to the assignor. */ public record GroupSpecImpl( Map members, - Map configs + AssignmentConfigs configs ) implements GroupSpec { public GroupSpecImpl { - // Both maps are exposed to a custom assignor through the public GroupSpec interface. + // The map is exposed to a custom assignor through the public GroupSpec interface. members = Collections.unmodifiableMap(Objects.requireNonNull(members)); - configs = Collections.unmodifiableMap(Objects.requireNonNull(configs)); + Objects.requireNonNull(configs); } @Override 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 529b07fd4fd74..1d735764b0d67 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 @@ -95,9 +95,7 @@ private static LinkedList taskIds(final TopologyDescriber topologyDescri private static LocalState initialize(final GroupSpec groupSpec, final TopologyDescriber topologyDescriber) { final LocalState localState = new LocalState(); - localState.numStandbyReplicas = - groupSpec.configs().isEmpty() ? 0 - : Integer.parseInt(groupSpec.configs().get("num.standby.replicas")); + localState.numStandbyReplicas = groupSpec.configs().numStandbyReplicas(); // Helpers for computing active tasks per member, and tasks per member localState.totalActiveTasks = 0; 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 c1bbff9439841..c0b58a3924e59 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 @@ -154,6 +154,7 @@ 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.AssignmentConfigsImpl; import org.apache.kafka.image.MetadataDelta; import org.apache.kafka.image.MetadataImage; import org.apache.kafka.image.MetadataProvenance; @@ -23654,7 +23655,7 @@ public void testStreamsGroupDynamicConfigs() { .setWarmupTasks(List.of())); assertEquals(2, result.response().data().memberEpoch()); assertEquals( - getDefaultAssignmentConfigs(), + new AssignmentConfigsImpl(GroupCoordinatorConfig.STREAMS_GROUP_NUM_STANDBY_REPLICAS_DEFAULT, List.of()), assignor.lastPassedAssignmentConfigs() ); @@ -23693,7 +23694,7 @@ public void testStreamsGroupDynamicConfigs() { // Verify that the new number of standby replicas is used assertEquals( - Map.of("num.standby.replicas", "2"), + new AssignmentConfigsImpl(2, List.of()), assignor.lastPassedAssignmentConfigs() ); @@ -24041,7 +24042,7 @@ public void testStreamsGroupEvaluatedConfigs() { context.assertSessionTimeout(groupId, memberId, GroupCoordinatorConfig.STREAMS_GROUP_SESSION_TIMEOUT_MS_DEFAULT); assertEquals( - getDefaultAssignmentConfigs(), + new AssignmentConfigsImpl(GroupCoordinatorConfig.STREAMS_GROUP_NUM_STANDBY_REPLICAS_DEFAULT, List.of()), assignor.lastPassedAssignmentConfigs()); assertEquals(GroupCoordinatorConfig.STREAMS_GROUP_TASK_OFFSET_INTERVAL_MS_DEFAULT, result.response().data().taskOffsetIntervalMs()); @@ -24081,7 +24082,7 @@ public void testStreamsGroupEvaluatedConfigs() { // Verify that the number of standby replicas is evaluated to max, // and task offset interval is evaluated to min assertEquals( - Map.of("num.standby.replicas", String.valueOf(GroupCoordinatorConfig.STREAMS_GROUP_MAX_STANDBY_REPLICAS_DEFAULT)), + new AssignmentConfigsImpl(GroupCoordinatorConfig.STREAMS_GROUP_MAX_STANDBY_REPLICAS_DEFAULT, List.of()), assignor.lastPassedAssignmentConfigs()); assertEquals(GroupCoordinatorConfig.STREAMS_GROUP_MIN_TASK_OFFSET_INTERVAL_MS_DEFAULT, result.response().data().taskOffsetIntervalMs()); 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 93bab5825e128..d1fc1de6188af 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,6 +16,7 @@ */ package org.apache.kafka.coordinator.group.streams; +import org.apache.kafka.coordinator.group.api.streams.assignor.AssignmentConfigs; 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; @@ -31,7 +32,7 @@ public class MockTaskAssignor implements TaskAssignor { private final String name; private GroupAssignment preparedGroupAssignment = null; - private Map assignmentConfigs = Map.of(); + private AssignmentConfigs assignmentConfigs = null; public MockTaskAssignor(String name) { this.name = name; @@ -53,7 +54,7 @@ public void prepareGroupAssignment(Map memberAssignments) { }))); } - public Map lastPassedAssignmentConfigs() { + public AssignmentConfigs lastPassedAssignmentConfigs() { return assignmentConfigs; } 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 e9f17f11a9324..a09001b24f491 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 @@ -27,6 +27,7 @@ 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.AssignmentConfigsImpl; import org.apache.kafka.coordinator.group.streams.assignor.GroupSpecImpl; import org.apache.kafka.coordinator.group.streams.assignor.MemberMetadataAndStateImpl; import org.apache.kafka.coordinator.group.streams.topics.ConfiguredSubtopology; @@ -479,7 +480,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(memberMetadataMap, new HashMap<>()); + GroupSpecImpl groupSpec = new GroupSpecImpl(memberMetadataMap, AssignmentConfigsImpl.DEFAULT); // 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/AssignmentConfigsImplTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/AssignmentConfigsImplTest.java new file mode 100644 index 0000000000000..2e619fa0451a5 --- /dev/null +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/AssignmentConfigsImplTest.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.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; + +public class AssignmentConfigsImplTest { + + @Test + void testFromEmptyMap() { + // A group metadata record written before the last assignment configs were persisted replays as an empty map. + assertEquals(AssignmentConfigsImpl.DEFAULT, AssignmentConfigsImpl.fromMap(Map.of())); + } + + @Test + void testFromMapWithoutRackAwareAssignmentTags() { + // The tags are only put in the map when any are configured. + assertEquals( + new AssignmentConfigsImpl(2, List.of()), + AssignmentConfigsImpl.fromMap(Map.of("num.standby.replicas", "2")) + ); + } + + @Test + void testFromMap() { + assertEquals( + new AssignmentConfigsImpl(1, List.of("tag1", "tag2")), + AssignmentConfigsImpl.fromMap(Map.of( + "num.standby.replicas", "1", + "rack.aware.assignment.tags", " tag1 , tag2 " + )) + ); + } + + @Test + void testRackAwareAssignmentTagsAreUnmodifiable() { + List tags = new ArrayList<>(List.of("tag1")); + AssignmentConfigsImpl configs = new AssignmentConfigsImpl(0, tags); + + tags.add("tag2"); + assertEquals(List.of("tag1"), configs.rackAwareAssignmentTags()); + assertThrows(UnsupportedOperationException.class, () -> configs.rackAwareAssignmentTags().add("tag2")); + } +} 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 976289b5e8e59..6d49a4f844913 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 @@ -19,13 +19,14 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import java.util.ArrayList; import java.util.HashMap; +import java.util.List; import java.util.Map; 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 { @@ -53,7 +54,7 @@ void setUp() { groupSpec = new GroupSpecImpl( members, - new HashMap<>() + new AssignmentConfigsImpl(2, new ArrayList<>(List.of("test-tag"))) ); } @@ -80,13 +81,14 @@ void testMemberNotFound() { @Test void testConfigs() { - assertTrue(groupSpec.configs().isEmpty()); + assertEquals(2, groupSpec.configs().numStandbyReplicas()); + assertEquals(List.of("test-tag"), groupSpec.configs().rackAwareAssignmentTags()); } @Test void testMembersAndConfigsAreUnmodifiable() { assertThrows(UnsupportedOperationException.class, () -> groupSpec.members().put("other-member", member)); - assertThrows(UnsupportedOperationException.class, () -> groupSpec.configs().put("key", "value")); + assertThrows(UnsupportedOperationException.class, () -> groupSpec.configs().rackAwareAssignmentTags().add("other-tag")); } } 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 a4d4d914aab51..0d93d3a5f6209 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 @@ -23,7 +23,6 @@ import org.junit.jupiter.api.Test; -import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; @@ -54,7 +53,7 @@ public void testZeroMembers() { TaskAssignorException ex = assertThrows(TaskAssignorException.class, () -> assignor.assign( new GroupSpecImpl( Map.of(), - new HashMap<>() + AssignmentConfigsImpl.DEFAULT ), new TopologyDescriberImpl(5, List.of("test-subtopology")) )); @@ -92,7 +91,7 @@ public void testDoubleAssignment() { TaskAssignorException ex = assertThrows(TaskAssignorException.class, () -> assignor.assign( new GroupSpecImpl( Map.of("member1", memberMetadata1, "member2", memberMetadata2), - new HashMap<>() + AssignmentConfigsImpl.DEFAULT ), new TopologyDescriberImpl(5, List.of("test-subtopology")) )); @@ -106,7 +105,7 @@ public void testBasicScenario() { final GroupAssignment result = assignor.assign( new GroupSpecImpl( Map.of(), - new HashMap<>() + AssignmentConfigsImpl.DEFAULT ), new TopologyDescriberImpl(5, List.of()) ); @@ -133,7 +132,7 @@ public void testSingleMember() { final GroupAssignment result = assignor.assign( new GroupSpecImpl( Map.of("test_member", memberMetadata), - new HashMap<>() + AssignmentConfigsImpl.DEFAULT ), new TopologyDescriberImpl(4, List.of("test-subtopology")) ); @@ -177,7 +176,7 @@ public void testTwoMembersTwoSubtopologies() { final GroupAssignment result = assignor.assign( new GroupSpecImpl( mkMap(mkEntry("test_member1", memberMetadata1), mkEntry("test_member2", memberMetadata2)), - new HashMap<>() + AssignmentConfigsImpl.DEFAULT ), new TopologyDescriberImpl(4, List.of("test-subtopology1", "test-subtopology2")) ); @@ -235,7 +234,7 @@ public void testTwoMembersTwoSubtopologiesStickiness() { final GroupAssignment result = assignor.assign( new GroupSpecImpl( mkMap(mkEntry("test_member1", memberMetadata1), mkEntry("test_member2", memberMetadata2)), - new HashMap<>() + AssignmentConfigsImpl.DEFAULT ), 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 f623f196f2337..412a64daadc61 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 @@ -46,7 +46,6 @@ public class StickyTaskAssignorTest { - public static final String NUM_STANDBY_REPLICAS_CONFIG = "num.standby.replicas"; private final StickyTaskAssignor assignor = new StickyTaskAssignor(); @Test @@ -64,7 +63,7 @@ public void shouldAssignOneActiveTaskToEachProcessWhenTaskCountSameAsProcessCoun final GroupAssignment result = assignor.assign( new GroupSpecImpl( mkMap(mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)), - new HashMap<>() + AssignmentConfigsImpl.DEFAULT ), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) ); @@ -93,7 +92,7 @@ public void shouldAssignTopicGroupIdEvenlyAcrossClientsWithNoStandByTasks() { mkEntry("member3_1", memberMetadata31), mkEntry("member3_2", memberMetadata32)); final GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, Arrays.asList("test-subtopology1", "test-subtopology2")) ); @@ -122,7 +121,7 @@ public void shouldAssignTopicGroupIdEvenlyAcrossClientsWithStandByTasks() { final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, - mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, "1"))), + new AssignmentConfigsImpl(1, List.of())), new TopologyDescriberImpl(3, true, Arrays.asList("test-subtopology1", "test-subtopology2")) ); @@ -144,7 +143,7 @@ public void shouldNotMigrateActiveTaskToOtherProcess() { Map members = mkMap(mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) ); @@ -163,7 +162,7 @@ public void shouldNotMigrateActiveTaskToOtherProcess() { members = mkMap(mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) ); @@ -186,7 +185,7 @@ public void shouldMigrateActiveTasksToNewProcessWithoutChangingAllAssignments() mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) ); @@ -214,7 +213,7 @@ public void shouldAssignBasedOnCapacity() { mkEntry("member1", memberMetadata1), mkEntry("member2_1", memberMetadata21), mkEntry("member2_2", memberMetadata22)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) ); @@ -240,7 +239,7 @@ public void shouldAssignTasksEvenlyWithUnequalTopicGroupSizes() { mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl2() ); @@ -271,7 +270,7 @@ public void shouldKeepActiveTaskStickinessWhenMoreClientThanActiveTasks() { mkEntry("member3", memberMetadata3), mkEntry("member4", memberMetadata4), mkEntry("member5", memberMetadata5)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) ); @@ -305,7 +304,7 @@ public void shouldKeepActiveTaskStickinessWhenMoreClientThanActiveTasks() { mkEntry("member3", memberMetadata3), mkEntry("member4", memberMetadata4), mkEntry("member5", memberMetadata5)); result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) ); @@ -338,7 +337,7 @@ public void shouldAssignTasksToClientWithPreviousStandbyTasks() { mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) ); @@ -364,7 +363,7 @@ public void shouldNotAssignStandbyTasksToClientWithPreviousStandbyTasksAndCurren mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, "1"))), + new GroupSpecImpl(members, new AssignmentConfigsImpl(1, List.of())), new TopologyDescriberImpl(2, true, List.of("test-subtopology")) ); @@ -397,7 +396,7 @@ public void shouldAssignBasedOnCapacityWhenMultipleClientHaveStandbyTasks() { mkEntry("member2_1", memberMetadata21), mkEntry("member2_2", memberMetadata22)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) ); @@ -428,7 +427,7 @@ public void shouldAssignStandbyTasksToDifferentClientThanCorrespondingActiveTask final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, - mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, "1"))), + new AssignmentConfigsImpl(1, List.of())), new TopologyDescriberImpl(4, true, List.of("test-subtopology")) ); @@ -464,7 +463,7 @@ public void shouldAssignMultipleReplicasOfStandbyTask() { final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, - mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, "2"))), + new AssignmentConfigsImpl(2, List.of())), new TopologyDescriberImpl(3, true, List.of("test-subtopology")) ); @@ -481,7 +480,7 @@ public void shouldNotAssignStandbyTaskReplicasWhenNoClientAvailableWithoutHaving final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, - mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, "1"))), + new AssignmentConfigsImpl(1, List.of())), new TopologyDescriberImpl(1, true, List.of("test-subtopology")) ); @@ -499,7 +498,7 @@ public void shouldAssignActiveAndStandbyTasks() { final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, - mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, "1"))), + new AssignmentConfigsImpl(1, List.of())), new TopologyDescriberImpl(3, true, List.of("test-subtopology")) ); @@ -519,7 +518,7 @@ public void shouldAssignAtLeastOneTaskToEachClientIfPossible() { mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) ); @@ -541,7 +540,7 @@ public void shouldAssignEachActiveTaskToOneClientWhenMoreClientsThanTasks() { mkEntry("member4", memberMetadata4), mkEntry("member5", memberMetadata5), mkEntry("member6", memberMetadata6)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) ); @@ -563,7 +562,7 @@ public void shouldBalanceActiveAndStandbyTasksAcrossAvailableClients() { final GroupAssignment result = assignor.assign( new GroupSpecImpl(members, - mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, "1"))), + new AssignmentConfigsImpl(1, List.of())), new TopologyDescriberImpl(3, true, List.of("test-subtopology")) ); @@ -581,7 +580,7 @@ public void shouldAssignMoreTasksToClientWithMoreCapacity() { mkEntry("member1", memberMetadata1), mkEntry("member2_1", memberMetadata21), mkEntry("member2_2", memberMetadata22)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, Arrays.asList("test-subtopology0", "test-subtopology1", "test-subtopology2", "test-subtopology3")) ); @@ -599,7 +598,7 @@ public void shouldReBalanceTasksAcrossAllClientsWhenCapacityAndTaskCountTheSame( mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3), mkEntry("member4", memberMetadata4)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(4, false, List.of("test-subtopology")) ); @@ -618,7 +617,7 @@ public void shouldReBalanceTasksAcrossClientsWhenCapacityLessThanTaskCount() { mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(4, false, List.of("test-subtopology")) ); @@ -636,7 +635,7 @@ public void shouldRebalanceTasksToClientsBasedOnCapacity() { mkEntry("member2", memberMetadata2), mkEntry("member3_1", memberMetadata31), mkEntry("member3_2", memberMetadata32)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(3, false, List.of("test-subtopology")) ); @@ -655,7 +654,7 @@ public void shouldMoveMinimalNumberOfTasksWhenPreviouslyAboveCapacityAndNewClien mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(4, false, List.of("test-subtopology")) ); @@ -676,7 +675,7 @@ public void shouldNotMoveAnyTasksWhenNewTasksAdded() { mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(6, false, List.of("test-subtopology")) ); @@ -697,7 +696,7 @@ public void shouldAssignNewTasksToNewClientWhenPreviousTasksAssignedToOldClients mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(6, false, List.of("test-subtopology")) ); @@ -730,7 +729,7 @@ public void shouldAssignTasksNotPreviouslyActiveToNewClient() { mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3), mkEntry("newMember", newMemberSpec)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(4, false, Arrays.asList("test-subtopology0", "test-subtopology1", "test-subtopology2")) ); @@ -762,7 +761,7 @@ public void shouldAssignTasksNotPreviouslyActiveToMultipleNewClients() { mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("bounce_member1", bounce1), mkEntry("bounce_member2", bounce2)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(4, false, Arrays.asList("test-subtopology0", "test-subtopology1", "test-subtopology2")) ); @@ -784,7 +783,7 @@ public void shouldAssignTasksToNewClient() { mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(2, false, List.of("test-subtopology")) ); @@ -800,7 +799,7 @@ public void shouldAssignTasksToNewClientWithoutFlippingAssignmentBetweenExisting mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("newMember", newMemberSpec)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(6, false, List.of("test-subtopology")) ); @@ -826,7 +825,7 @@ public void shouldAssignTasksToNewClientWithoutFlippingAssignmentBetweenExisting mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("newMember", newMemberSpec)); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, new HashMap<>()), + new GroupSpecImpl(members, AssignmentConfigsImpl.DEFAULT), new TopologyDescriberImpl(7, false, List.of("test-subtopology")) ); @@ -855,7 +854,7 @@ public void shouldHandleLargeNumberOfTasksWithStandbyAssignment() { } GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, String.valueOf(numStandbyReplicas)))), + new GroupSpecImpl(members, new AssignmentConfigsImpl(numStandbyReplicas, List.of())), new TopologyDescriberImpl(numTasks, true, List.of("test-subtopology")) ); @@ -926,7 +925,7 @@ public void shouldHandleOddNumberOfClientsWithStandbyTasks() { } GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, String.valueOf(numStandbyReplicas)))), + new GroupSpecImpl(members, new AssignmentConfigsImpl(numStandbyReplicas, List.of())), new TopologyDescriberImpl(numTasks, true, List.of("test-subtopology")) ); @@ -977,7 +976,7 @@ public void shouldHandleHighStandbyReplicaCount() { } GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, String.valueOf(numStandbyReplicas)))), + new GroupSpecImpl(members, new AssignmentConfigsImpl(numStandbyReplicas, List.of())), new TopologyDescriberImpl(numTasks, true, List.of("test-subtopology")) ); @@ -1021,7 +1020,7 @@ public void shouldHandleLargeNumberOfSubtopologiesWithStandbyTasks() { } GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, String.valueOf(numStandbyReplicas)))), + new GroupSpecImpl(members, new AssignmentConfigsImpl(numStandbyReplicas, List.of())), new TopologyDescriberImpl(5, true, subtopologies) // 5 tasks per subtopology ); @@ -1053,7 +1052,7 @@ public void shouldHandleEdgeCaseWithSingleClientAndMultipleStandbyReplicas() { ); GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, String.valueOf(numStandbyReplicas)))), + new GroupSpecImpl(members, new AssignmentConfigsImpl(numStandbyReplicas, List.of())), new TopologyDescriberImpl(numTasks, true, List.of("test-subtopology")) ); @@ -1078,7 +1077,7 @@ public void shouldHandleEdgeCaseWithMoreStandbyReplicasThanAvailableClients() { } GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, String.valueOf(numStandbyReplicas)))), + new GroupSpecImpl(members, new AssignmentConfigsImpl(numStandbyReplicas, List.of())), new TopologyDescriberImpl(numTasks, true, List.of("test-subtopology")) ); @@ -1120,7 +1119,7 @@ public void shouldReassignTasksWhenNewNodeJoinsWithExistingActiveAndStandbyAssig mkEntry("member1", memberMetadata1), mkEntry("member2", memberMetadata2), mkEntry("member3", memberMetadata3)); final GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, "1"))), + new GroupSpecImpl(members, new AssignmentConfigsImpl(1, List.of())), new TopologyDescriberImpl(4, true, List.of("test-subtopology")) ); @@ -1161,7 +1160,7 @@ public void shouldRangeAssignTasksWhenScalingUp() { // Two subtopologies with 2 tasks each (4 tasks total) with standby replicas enabled final GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, String.valueOf(1)))), + new GroupSpecImpl(members, new AssignmentConfigsImpl(1, List.of())), new TopologyDescriberImpl(2, true, Arrays.asList("test-subtopology1", "test-subtopology2")) ); @@ -1207,7 +1206,7 @@ public void shouldRangeAssignTasksWhenStartingEmpty() { // Two subtopologies with 2 tasks each (4 tasks total) with standby replicas enabled final GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, String.valueOf(1)))), + new GroupSpecImpl(members, new AssignmentConfigsImpl(1, List.of())), new TopologyDescriberImpl(2, true, Arrays.asList("test-subtopology1", "test-subtopology2")) ); @@ -1277,7 +1276,7 @@ public void shouldAssignStandbyTaskToPreviousOwnerBasedOnBelowQuotaCondition() { // We have 2 active tasks + 1 standby replica = 4 total tasks // Quota per process = 4 tasks / 3 processes = 1.33 -> 2 tasks per process final GroupAssignment result = assignor.assign( - new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, "1"))), + new GroupSpecImpl(members, new AssignmentConfigsImpl(1, List.of())), new TopologyDescriberImpl(2, true, List.of("test-subtopology")) ); 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 b1478498a3f42..caa2a619e54ef 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,6 +16,7 @@ */ package org.apache.kafka.jmh.assignor; +import org.apache.kafka.coordinator.group.api.streams.assignor.AssignmentConfigs; 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.GroupSpecImpl; @@ -44,7 +45,7 @@ public class StreamsAssignorBenchmarkUtils { */ public static GroupSpec createGroupSpec( Map members, - Map assignmentConfigs + AssignmentConfigs assignmentConfigs ) { Map memberSpecs = new HashMap<>(); 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 b25c441cf821e..d5f85172c4e56 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,6 +17,7 @@ package org.apache.kafka.jmh.assignor; import org.apache.kafka.coordinator.common.runtime.CoordinatorMetadataImage; +import org.apache.kafka.coordinator.group.api.streams.assignor.AssignmentConfigs; 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; @@ -24,6 +25,7 @@ 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.AssignmentConfigsImpl; import org.apache.kafka.coordinator.group.streams.assignor.GroupSpecImpl; import org.apache.kafka.coordinator.group.streams.assignor.MemberMetadataAndStateImpl; import org.apache.kafka.coordinator.group.streams.assignor.StickyTaskAssignor; @@ -91,7 +93,7 @@ public enum AssignmentType { private TopologyDescriber topologyDescriber; - private Map assignmentConfigs; + private AssignmentConfigs assignmentConfigs; @Setup(Level.Trial) public void setup() { @@ -106,10 +108,7 @@ public void setup() { taskAssignor = new StickyTaskAssignor(); Map members = createMembers(); - this.assignmentConfigs = Map.of( - "num.standby.replicas", - Integer.toString(standbyReplicas) - ); + this.assignmentConfigs = new AssignmentConfigsImpl(standbyReplicas, List.of()); this.groupSpec = StreamsAssignorBenchmarkUtils.createGroupSpec(members, assignmentConfigs); if (assignmentType == AssignmentType.INCREMENTAL) {