From 2ae93771eba6aad748fdf16242253698ec9beab2 Mon Sep 17 00:00:00 2001 From: PoAn Yang Date: Mon, 17 Nov 2025 21:56:28 +0800 Subject: [PATCH 1/6] KAFKA-19387: Add rack awareness assignment to UniformHomogeneousAssignmentBuilder Signed-off-by: PoAn Yang --- .../UniformHomogeneousAssignmentBuilder.java | 163 ++++++- .../coordinator/group/assignor/Utils.java | 42 ++ ...OptimizedUniformAssignmentBuilderTest.java | 425 ++++++++++++++++++ .../coordinator/group/assignor/UtilsTest.java | 101 +++++ 4 files changed, 711 insertions(+), 20 deletions(-) create mode 100644 group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/Utils.java create mode 100644 group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/UtilsTest.java diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java index d3a5c35a86165..27ae8ce147b39 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java @@ -26,10 +26,14 @@ import org.apache.kafka.server.common.TopicIdPartition; import java.util.ArrayList; +import java.util.Collection; import java.util.HashMap; import java.util.HashSet; +import java.util.Iterator; +import java.util.LinkedList; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.Set; /** @@ -41,11 +45,13 @@ *
  • Balance: Ensure partitions are distributed equally among all members. * The difference in assignments sizes between any two members * should not exceed one partition.
  • + *
  • Rack Awareness: When feasible, aim to assign partitions to members + * located on the same rack thus avoiding cross-zone traffic.
  • *
  • Stickiness: Minimize partition movements among members by retaining * as much of the existing assignment as possible.
  • * * The assignment builder prioritizes the properties in the following order: - * Balance > Stickiness. + * Balance > Rack Awareness > Stickiness. */ public class UniformHomogeneousAssignmentBuilder { /** @@ -66,7 +72,7 @@ public class UniformHomogeneousAssignmentBuilder { /** * The members that are below their quota. */ - private final List unfilledMembers; + private final LinkedList unfilledMembers; /** * The partitions that still need to be assigned. @@ -85,6 +91,11 @@ public class UniformHomogeneousAssignmentBuilder { */ private int minimumMemberQuota; + /** + * Whether to use rack aware assignment strategy. + */ + private final boolean useRackStrategy; + /** * The number of members to receive an extra partition beyond the minimum quota. * Example: If there are 11 partitions to be distributed among 3 members, @@ -97,10 +108,34 @@ public class UniformHomogeneousAssignmentBuilder { this.subscribedTopicDescriber = subscribedTopicDescriber; this.subscribedTopicIds = new HashSet<>(groupSpec.memberSubscription(groupSpec.memberIds().iterator().next()) .subscribedTopicIds()); - this.unfilledMembers = new ArrayList<>(); - this.unassignedPartitions = new ArrayList<>(); + this.unfilledMembers = new LinkedList<>(); + this.unassignedPartitions = new LinkedList<>(); this.targetAssignment = new HashMap<>(); + + Set allMemberRacks = groupSpec.memberIds().stream() + .map(memberId -> groupSpec.memberSubscription(memberId).rackId()) + .filter(Optional::isPresent) + .map(Optional::get) + .collect(java.util.stream.Collectors.toSet()); + Map> racksPerPartition = this.subscribedTopicIds.stream() + .flatMap(topicId -> { + int partitionCount = subscribedTopicDescriber.numPartitions(topicId); + List topicPartitions = new ArrayList<>(); + for (int partitionId = 0; partitionId < partitionCount; partitionId++) { + topicPartitions.add(new TopicIdPartition(topicId, partitionId)); + } + return topicPartitions.stream(); + }) + .collect(java.util.stream.Collectors.toMap( + tip -> tip, + tip -> subscribedTopicDescriber.racksForPartition(tip.topicId(), tip.partitionId()) + )); + Set allPartitionRacks = racksPerPartition.values().stream() + .flatMap(Collection::stream) + .collect(java.util.stream.Collectors.toSet()); + + this.useRackStrategy = Utils.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition); } /** @@ -139,6 +174,11 @@ public GroupAssignment build() throws PartitionAssignorException { // exceed the maximum quota assigned to each member. maybeRevokePartitions(); + // Assign the unassigned partitions to the members if rack matches. + if (useRackStrategy) { + assignRackAwarenessRemainingPartitions(); + } + // Assign the unassigned partitions to the members with space. assignRemainingPartitions(); @@ -154,11 +194,7 @@ public GroupAssignment build() throws PartitionAssignorException { */ @SuppressWarnings({"CyclomaticComplexity", "NPathComplexity"}) private void maybeRevokePartitions() { - int memberCount = groupSpec.memberIds().size(); - int memberIndex = -1; for (String memberId : groupSpec.memberIds()) { - memberIndex++; - Map> oldAssignment = groupSpec.memberAssignment(memberId).partitions(); Map> newAssignment = null; @@ -180,11 +216,17 @@ private void maybeRevokePartitions() { Set partitions = topicPartitions.getValue(); if (subscribedTopicIds.contains(topicId)) { - if (partitions.size() <= quota) { + if (partitions.size() <= quota && !useRackStrategy) { quota -= partitions.size(); } else { for (Integer partition : partitions) { - if (quota > 0) { + // Keep the partition when: + // 1. We still have quota left to keep it. + // 2-1. Don't use rack strategy, so we can keep it. + // 2-2. Use rack strategy and the member rack matches the partition racks. + if (quota > 0 && (!useRackStrategy || Utils.isRackMatch( + groupSpec.memberSubscription(memberId).rackId(), + subscribedTopicDescriber.racksForPartition(topicId, partition)))) { quota--; } else { if (newAssignment == null) { @@ -215,21 +257,21 @@ private void maybeRevokePartitions() { } if (quota > 0 && - quotaHasExtraPartition && - memberCount - memberIndex > remainingMembersToGetAnExtraPartition) { - // Give up the extra partition quota for another member to claim, - // unless this member is one of the last remainingMembersToGetAnExtraPartition - // members in the list and must take the extra partition. + quotaHasExtraPartition) { + // Give up the extra partition quota for another member to claim. quota--; quotaHasExtraPartition = false; } - if (quota > 0) { - unfilledMembers.add(new MemberWithRemainingQuota(memberId, quota)); + if (quota == 0 && quotaHasExtraPartition) { + remainingMembersToGetAnExtraPartition--; } - if (quotaHasExtraPartition) { - remainingMembersToGetAnExtraPartition--; + // An unfilled member is a member that can still take more partitions. + // 1. It has quota left. + // 2. It doesn't have quota left, but there are remaining extra partition, and it hasn't claimed one yet. + if (quota > 0 || (remainingMembersToGetAnExtraPartition > 0 && !quotaHasExtraPartition)) { + unfilledMembers.add(new MemberWithRemainingQuota(memberId, quota)); } if (newAssignment == null) { @@ -240,12 +282,74 @@ private void maybeRevokePartitions() { } } + /** + * Assign the unassigned partitions to the unfilled members if member and partition racks are matched. + */ + private void assignRackAwarenessRemainingPartitions() { + for (var unassignedPartitionsIter = unassignedPartitions.iterator(); unassignedPartitionsIter.hasNext(); ) { + TopicIdPartition tip = unassignedPartitionsIter.next(); + boolean isPartitionAssigned = false; + + for (var unfilledMembersIter = unfilledMembers.iterator(); unfilledMembersIter.hasNext(); ) { + MemberWithRemainingQuota unfilledMember = unfilledMembersIter.next(); + if (unfilledMember.remainingQuota() == 0 && remainingMembersToGetAnExtraPartition == 0) { + unfilledMembersIter.remove(); + continue; + } + + String memberId = unfilledMember.memberId; + if (!Utils.isRackMatch(groupSpec.memberSubscription(memberId).rackId(), + subscribedTopicDescriber.racksForPartition(tip.topicId(), tip.partitionId()))) { + continue; + } + + Map> newAssignment = targetAssignment.get(memberId).partitions(); + if (AssignorHelpers.isImmutableMap(newAssignment)) { + // If the new assignment is immutable, we must create a deep copy of it + // before altering it. + newAssignment = AssignorHelpers.deepCopyAssignment(newAssignment); + targetAssignment.put(memberId, new MemberAssignmentImpl(newAssignment)); + } + newAssignment + .computeIfAbsent(tip.topicId(), __ -> new HashSet<>()) + .add(tip.partitionId()); + + if (unfilledMember.remainingQuota() > 0) { + unfilledMember.decrementQuota(); + } else { + unfilledMembersIter.remove(); + remainingMembersToGetAnExtraPartition--; + } + + isPartitionAssigned = true; + break; + } + + if (isPartitionAssigned) { + unassignedPartitionsIter.remove(); + } + } + } + /** * Assign the unassigned partitions to the unfilled members. */ private void assignRemainingPartitions() { int unassignedPartitionIndex = 0; + // If member is one of the last remainingMembersToGetAnExtraPartition, it must take the extra partition. + for (Iterator it = unfilledMembers.descendingIterator(); it.hasNext(); ) { + MemberWithRemainingQuota unfilledMember = it.next(); + if (remainingMembersToGetAnExtraPartition > 0) { + remainingMembersToGetAnExtraPartition--; + unfilledMember.incrementQuota(); + } + + if (unfilledMember.remainingQuota() == 0) { + it.remove(); + } + } + for (MemberWithRemainingQuota unfilledMember : unfilledMembers) { String memberId = unfilledMember.memberId; int remainingQuota = unfilledMember.remainingQuota; @@ -272,6 +376,25 @@ private void assignRemainingPartitions() { } } - private record MemberWithRemainingQuota(String memberId, int remainingQuota) { + private static class MemberWithRemainingQuota { + private final String memberId; + private int remainingQuota; + + MemberWithRemainingQuota(String memberId, int remainingQuota) { + this.memberId = memberId; + this.remainingQuota = remainingQuota; + } + + int remainingQuota() { + return remainingQuota; + } + + void incrementQuota() { + remainingQuota++; + } + + void decrementQuota() { + remainingQuota--; + } } } diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/Utils.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/Utils.java new file mode 100644 index 0000000000000..cf47125528d7f --- /dev/null +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/Utils.java @@ -0,0 +1,42 @@ +/* + * 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.assignor; + +import org.apache.kafka.server.common.TopicIdPartition; + +import java.util.Collections; +import java.util.Map; +import java.util.Optional; +import java.util.Set; + +public class Utils { + public static boolean isRackMatch(Optional memberRackId, Set partitionRackIds) { + return memberRackId.isPresent() && partitionRackIds.contains(memberRackId.get()); + } + + public static boolean useRackAwareAssignment( + Set allMemberRacks, + Set allPartitionRacks, + Map> racksPerPartition + ) { + if (allMemberRacks.isEmpty() || Collections.disjoint(allMemberRacks, allPartitionRacks)) + return false; + else { + return !racksPerPartition.values().stream().allMatch(allPartitionRacks::equals); + } + } +} \ No newline at end of file diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/OptimizedUniformAssignmentBuilderTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/OptimizedUniformAssignmentBuilderTest.java index 284bde18dfc24..435556fda579c 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/OptimizedUniformAssignmentBuilderTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/OptimizedUniformAssignmentBuilderTest.java @@ -17,6 +17,7 @@ package org.apache.kafka.coordinator.group.assignor; import org.apache.kafka.common.Uuid; +import org.apache.kafka.coordinator.common.runtime.CoordinatorMetadataImage; import org.apache.kafka.coordinator.common.runtime.KRaftCoordinatorMetadataImage; import org.apache.kafka.coordinator.common.runtime.MetadataImageBuilder; import org.apache.kafka.coordinator.group.api.assignor.GroupAssignment; @@ -632,6 +633,430 @@ public void testReassignmentStickinessWhenAlreadyBalanced() { checkValidityAndBalance(members, computedAssignment); } + + @Test + public void testFirstAssignmentThreeMembersOneTopicWithRacks() { + CoordinatorMetadataImage coordinatorMetadataImage = new MetadataImageBuilder() + .addTopic(topic1Uuid, topic1Name, 6) + .addRacks() + .buildCoordinatorMetadataImage(); + + Map members = new TreeMap<>(); + + members.put(memberA, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack0"), + Optional.empty(), + Set.of(topic1Uuid), + Assignment.EMPTY + )); + + members.put(memberB, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack1"), + Optional.empty(), + Set.of(topic1Uuid), + Assignment.EMPTY + )); + + members.put(memberC, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack2"), + Optional.empty(), + Set.of(topic1Uuid), + Assignment.EMPTY + )); + + GroupSpec groupSpec = new GroupSpecImpl( + members, + HOMOGENEOUS, + Map.of() + ); + SubscribedTopicDescriberImpl subscribedTopicMetadata = new SubscribedTopicDescriberImpl( + coordinatorMetadataImage + ); + + GroupAssignment computedAssignment = assignor.assign( + groupSpec, + subscribedTopicMetadata + ); + + Map>> expectedAssignment = new HashMap<>(); + expectedAssignment.put(memberA, mkAssignment( + mkTopicAssignment(topic1Uuid, 0, 3) + )); + expectedAssignment.put(memberB, mkAssignment( + mkTopicAssignment(topic1Uuid, 1, 4) + )); + expectedAssignment.put(memberC, mkAssignment( + mkTopicAssignment(topic1Uuid, 2, 5) + )); + + assertAssignment(expectedAssignment, computedAssignment); + checkValidityAndBalance(members, computedAssignment); + } + + @Test + public void testFirstAssignmentThreeMembersOneTopicWithAllRacksMismatched() { + CoordinatorMetadataImage coordinatorMetadataImage = new MetadataImageBuilder() + .addTopic(topic1Uuid, topic1Name, 6) + .addRacks() + .buildCoordinatorMetadataImage(); + + Map members = new TreeMap<>(); + + members.put(memberA, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack5"), + Optional.empty(), + Set.of(topic1Uuid), + Assignment.EMPTY + )); + + members.put(memberB, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack6"), + Optional.empty(), + Set.of(topic1Uuid), + Assignment.EMPTY + )); + + members.put(memberC, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack7"), + Optional.empty(), + Set.of(topic1Uuid), + Assignment.EMPTY + )); + + GroupSpec groupSpec = new GroupSpecImpl( + members, + HOMOGENEOUS, + Map.of() + ); + SubscribedTopicDescriberImpl subscribedTopicMetadata = new SubscribedTopicDescriberImpl( + coordinatorMetadataImage + ); + + GroupAssignment computedAssignment = assignor.assign( + groupSpec, + subscribedTopicMetadata + ); + + Map>> expectedAssignment = new HashMap<>(); + expectedAssignment.put(memberA, mkAssignment( + mkTopicAssignment(topic1Uuid, 0, 1) + )); + expectedAssignment.put(memberB, mkAssignment( + mkTopicAssignment(topic1Uuid, 2, 3) + )); + expectedAssignment.put(memberC, mkAssignment( + mkTopicAssignment(topic1Uuid, 4, 5) + )); + + assertAssignment(expectedAssignment, computedAssignment); + checkValidityAndBalance(members, computedAssignment); + } + + @Test + public void testFirstAssignmentTwoMembersOneTopicWithMembersInSameRack() { + CoordinatorMetadataImage coordinatorMetadataImage = new MetadataImageBuilder() + .addTopic(topic1Uuid, topic1Name, 4) + .addRacks() + .buildCoordinatorMetadataImage(); + + Map members = new TreeMap<>(); + + members.put(memberA, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack0"), + Optional.empty(), + Set.of(topic1Uuid), + Assignment.EMPTY + )); + + members.put(memberB, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack0"), + Optional.empty(), + Set.of(topic1Uuid), + Assignment.EMPTY + )); + + GroupSpec groupSpec = new GroupSpecImpl( + members, + HOMOGENEOUS, + Map.of() + ); + SubscribedTopicDescriberImpl subscribedTopicMetadata = new SubscribedTopicDescriberImpl( + coordinatorMetadataImage + ); + + GroupAssignment computedAssignment = assignor.assign( + groupSpec, + subscribedTopicMetadata + ); + + Map>> expectedAssignment = new HashMap<>(); + expectedAssignment.put(memberA, mkAssignment( + mkTopicAssignment(topic1Uuid, 0, 3) + )); + expectedAssignment.put(memberB, mkAssignment( + mkTopicAssignment(topic1Uuid, 1, 2) + )); + + assertAssignment(expectedAssignment, computedAssignment); + checkValidityAndBalance(members, computedAssignment); + } + + @Test + public void testFirstAssignmentThreeMembersTwoTopicsWithOneMemberNoRack() { + MetadataImage metadataImage = new MetadataImageBuilder() + .addTopic(topic1Uuid, topic1Name, 6) + .addTopic(topic2Uuid, topic2Name, 3) + .addRacks() + .build(); + + Map members = new TreeMap<>(); + + members.put(memberA, new MemberSubscriptionAndAssignmentImpl( + Optional.empty(), + Optional.empty(), + Set.of(topic1Uuid, topic2Uuid), + Assignment.EMPTY + )); + + members.put(memberB, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack1"), + Optional.empty(), + Set.of(topic1Uuid, topic2Uuid), + Assignment.EMPTY + )); + + members.put(memberC, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack2"), + Optional.empty(), + Set.of(topic1Uuid, topic2Uuid), + Assignment.EMPTY + )); + + GroupSpec groupSpec = new GroupSpecImpl( + members, + HOMOGENEOUS, + Map.of() + ); + SubscribedTopicDescriberImpl subscribedTopicMetadata = new SubscribedTopicDescriberImpl( + new KRaftCoordinatorMetadataImage(metadataImage) + ); + + GroupAssignment computedAssignment = assignor.assign( + groupSpec, + subscribedTopicMetadata + ); + Map>> expectedAssignment = new HashMap<>(); + expectedAssignment.put(memberA, mkAssignment( + mkTopicAssignment(topic1Uuid, 3, 4, 5) + )); + expectedAssignment.put(memberB, mkAssignment( + mkTopicAssignment(topic1Uuid, 0), + mkTopicAssignment(topic2Uuid, 0, 1) + )); + expectedAssignment.put(memberC, mkAssignment( + mkTopicAssignment(topic1Uuid, 1, 2), + mkTopicAssignment(topic2Uuid, 2) + )); + + assertAssignment(expectedAssignment, computedAssignment); + checkValidityAndBalance(members, computedAssignment); + } + + @Test + public void testReassignmentWhenRackMismatched() { + MetadataImage metadataImage = new MetadataImageBuilder() + .addTopic(topic1Uuid, topic1Name, 6) + .addRacks() + .build(); + + Map members = new TreeMap<>(); + + members.put(memberA, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack0"), + Optional.empty(), + Set.of(topic1Uuid), + new Assignment(mkAssignment( + mkTopicAssignment(topic1Uuid, 0, 1) + )) + )); + + members.put(memberB, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack1"), + Optional.empty(), + Set.of(topic1Uuid), + new Assignment(mkAssignment( + mkTopicAssignment(topic1Uuid, 2, 4) + )) + )); + + members.put(memberC, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack2"), + Optional.empty(), + Set.of(topic1Uuid), + new Assignment(mkAssignment( + mkTopicAssignment(topic1Uuid, 3, 5) + )) + )); + + GroupSpec groupSpec = new GroupSpecImpl( + members, + HOMOGENEOUS, + invertedTargetAssignment(members) + ); + SubscribedTopicDescriberImpl subscribedTopicMetadata = new SubscribedTopicDescriberImpl( + new KRaftCoordinatorMetadataImage(metadataImage) + ); + + GroupAssignment computedAssignment = assignor.assign( + groupSpec, + subscribedTopicMetadata + ); + Map>> expectedAssignment = new HashMap<>(); + expectedAssignment.put(memberA, mkAssignment( + mkTopicAssignment(topic1Uuid, 0, 3) + )); + expectedAssignment.put(memberB, mkAssignment( + mkTopicAssignment(topic1Uuid, 1, 4) + )); + expectedAssignment.put(memberC, mkAssignment( + mkTopicAssignment(topic1Uuid, 2, 5) + )); + + assertAssignment(expectedAssignment, computedAssignment); + checkValidityAndBalance(members, computedAssignment); + } + + @Test + public void testReassignmentWhenAddingNewTopicWithRacks() { + // Add topic2 as new topic + MetadataImage metadataImage = new MetadataImageBuilder() + .addTopic(topic1Uuid, topic1Name, 6) + .addTopic(topic2Uuid, topic2Name, 2) + .addRacks() + .build(); + + Map members = new TreeMap<>(); + + members.put(memberA, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack0"), + Optional.empty(), + Set.of(topic1Uuid, topic2Uuid), + new Assignment(mkAssignment( + mkTopicAssignment(topic1Uuid, 0, 3) + )) + )); + + members.put(memberB, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack1"), + Optional.empty(), + Set.of(topic1Uuid, topic2Uuid), + new Assignment(mkAssignment( + mkTopicAssignment(topic1Uuid, 1, 4) + )) + )); + + members.put(memberC, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack2"), + Optional.empty(), + Set.of(topic1Uuid, topic2Uuid), + new Assignment(mkAssignment( + mkTopicAssignment(topic1Uuid, 2, 5) + )) + )); + + GroupSpec groupSpec = new GroupSpecImpl( + members, + HOMOGENEOUS, + invertedTargetAssignment(members) + ); + SubscribedTopicDescriberImpl subscribedTopicMetadata = new SubscribedTopicDescriberImpl( + new KRaftCoordinatorMetadataImage(metadataImage) + ); + + GroupAssignment computedAssignment = assignor.assign( + groupSpec, + subscribedTopicMetadata + ); + Map>> expectedAssignment = new HashMap<>(); + expectedAssignment.put(memberA, mkAssignment( + mkTopicAssignment(topic1Uuid, 0, 3), + mkTopicAssignment(topic2Uuid, 0) + )); + expectedAssignment.put(memberB, mkAssignment( + mkTopicAssignment(topic1Uuid, 1, 4), + mkTopicAssignment(topic2Uuid, 1) + )); + expectedAssignment.put(memberC, mkAssignment( + mkTopicAssignment(topic1Uuid, 2, 5) + )); + + assertAssignment(expectedAssignment, computedAssignment); + checkValidityAndBalance(members, computedAssignment); + } + + @Test + public void testReassignmentWhenAddingNewMemberWithRacks() { + MetadataImage metadataImage = new MetadataImageBuilder() + .addTopic(topic1Uuid, topic1Name, 4) + .addRacks() + .build(); + + Map members = new TreeMap<>(); + + members.put(memberA, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack0"), + Optional.empty(), + Set.of(topic1Uuid), + new Assignment(mkAssignment( + mkTopicAssignment(topic1Uuid, 0, 3) + )) + )); + + members.put(memberB, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack1"), + Optional.empty(), + Set.of(topic1Uuid), + new Assignment(mkAssignment( + mkTopicAssignment(topic1Uuid, 1, 2) + )) + )); + + // Add new member C + members.put(memberC, new MemberSubscriptionAndAssignmentImpl( + Optional.of("rack2"), + Optional.empty(), + Set.of(topic1Uuid), + Assignment.EMPTY + )); + + GroupSpec groupSpec = new GroupSpecImpl( + members, + HOMOGENEOUS, + invertedTargetAssignment(members) + ); + SubscribedTopicDescriberImpl subscribedTopicMetadata = new SubscribedTopicDescriberImpl( + new KRaftCoordinatorMetadataImage(metadataImage) + ); + + GroupAssignment computedAssignment = assignor.assign( + groupSpec, + subscribedTopicMetadata + ); + Map>> expectedAssignment = new HashMap<>(); + expectedAssignment.put(memberA, mkAssignment( + mkTopicAssignment(topic1Uuid, 0, 3) + )); + expectedAssignment.put(memberB, mkAssignment( + mkTopicAssignment(topic1Uuid, 1) + )); + expectedAssignment.put(memberC, mkAssignment( + mkTopicAssignment(topic1Uuid, 2) + )); + + assertAssignment(expectedAssignment, computedAssignment); + checkValidityAndBalance(members, computedAssignment); + } + /** * Verifies that the given assignment is valid with respect to the given subscriptions. * Validity requirements: diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/UtilsTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/UtilsTest.java new file mode 100644 index 0000000000000..1f04074f22e33 --- /dev/null +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/UtilsTest.java @@ -0,0 +1,101 @@ +/* + * 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.assignor; + +import org.apache.kafka.common.Uuid; +import org.apache.kafka.server.common.TopicIdPartition; + +import org.junit.jupiter.api.Test; + +import java.util.Map; +import java.util.Optional; +import java.util.Set; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +public class UtilsTest { + private static final Uuid TOPIC_ID = Uuid.randomUuid(); + + @Test + void testIsRackMatchWithEmptyMemberRackId() { + assertFalse(Utils.isRackMatch(Optional.empty(), Set.of("rack1", "rack2"))); + } + + @Test + void testIsRackMatchWithMatchingRack() { + assertTrue(Utils.isRackMatch(Optional.of("rack1"), Set.of("rack1", "rack2"))); + } + + @Test + void testIsRackMatchWithNonMatchingRack() { + assertFalse(Utils.isRackMatch(Optional.of("rack3"), Set.of("rack1", "rack2"))); + } + + @Test + void testIsRackMatchWithEmptyPartitionRacks() { + assertFalse(Utils.isRackMatch(Optional.of("rack1"), Set.of())); + } + + @Test + void testUseRackAwareAssignmentWithEmptyMemberRacks() { + Set allMemberRacks = Set.of(); + Set allPartitionRacks = Set.of("rack1", "rack2"); + Map> racksPerPartition = Map.of( + new TopicIdPartition(TOPIC_ID, 0), Set.of("rack1"), + new TopicIdPartition(TOPIC_ID, 1), Set.of("rack2") + ); + + assertFalse(Utils.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition)); + } + + @Test + void testUseRackAwareAssignmentWithDisjointRacks() { + Set allMemberRacks = Set.of("rack1", "rack2"); + Set allPartitionRacks = Set.of("rack3", "rack4"); + Map> racksPerPartition = Map.of( + new TopicIdPartition(TOPIC_ID, 0), Set.of("rack3"), + new TopicIdPartition(TOPIC_ID, 1), Set.of("rack4") + ); + + assertFalse(Utils.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition)); + } + + @Test + void testUseRackAwareAssignmentWithAllPartitionsHavingSameRacks() { + Set allMemberRacks = Set.of("rack1", "rack2"); + Set allPartitionRacks = Set.of("rack1", "rack2"); + Map> racksPerPartition = Map.of( + new TopicIdPartition(TOPIC_ID, 0), Set.of("rack1", "rack2"), + new TopicIdPartition(TOPIC_ID, 1), Set.of("rack1", "rack2") + ); + + assertFalse(Utils.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition)); + } + + @Test + void testUseRackAwareAssignmentWithDifferentRacksPerPartition() { + Set allMemberRacks = Set.of("rack1", "rack2"); + Set allPartitionRacks = Set.of("rack1", "rack2"); + Map> racksPerPartition = Map.of( + new TopicIdPartition(TOPIC_ID, 0), Set.of("rack1"), + new TopicIdPartition(TOPIC_ID, 1), Set.of("rack2") + ); + + assertTrue(Utils.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition)); + } +} From 81aa4210f37ee9b5ec62fb3e47ae968125602626 Mon Sep 17 00:00:00 2001 From: PoAn Yang Date: Wed, 19 Nov 2025 23:15:58 +0800 Subject: [PATCH 2/6] Move Utils functions to AssignorHelpers Signed-off-by: PoAn Yang --- .../group/assignor/AssignorHelpers.java | 32 ++++++++++++++ .../UniformHomogeneousAssignmentBuilder.java | 6 +-- .../coordinator/group/assignor/Utils.java | 42 ------------------- ...tilsTest.java => AssignorHelpersTest.java} | 18 ++++---- 4 files changed, 44 insertions(+), 54 deletions(-) delete mode 100644 group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/Utils.java rename group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/{UtilsTest.java => AssignorHelpersTest.java} (78%) diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/AssignorHelpers.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/AssignorHelpers.java index a2156c1febaed..f50ef51b84b0a 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/AssignorHelpers.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/AssignorHelpers.java @@ -17,11 +17,13 @@ package org.apache.kafka.coordinator.group.assignor; import org.apache.kafka.common.Uuid; +import org.apache.kafka.server.common.TopicIdPartition; import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.Map; +import java.util.Optional; import java.util.Set; /** @@ -69,4 +71,34 @@ static HashMap newHashMap(int numMappings) { static HashSet newHashSet(int numElements) { return new HashSet<>((int) (((numElements + 1) / 0.75f) + 1)); } + + + /** + * Checks if the member's rack matches any of the partition's racks. + * @param memberRackId The member's rack id. + * @param partitionRackIds The partition's rack ids. + * @return True if the member's rack matches any of the partition's racks, false otherwise. + */ + public static boolean isRackMatch(Optional memberRackId, Set partitionRackIds) { + return memberRackId.isPresent() && partitionRackIds.contains(memberRackId.get()); + } + + /** + * Determines whether rack-aware assignment should be used based on the provided racks. + * @param allMemberRacks The set of all member racks. + * @param allPartitionRacks The set of all partition racks. + * @param racksPerPartition A map of partitions to their respective racks. + * @return True if member racks and partition racks overlap and not all partitions have the same set of racks, false otherwise. + */ + public static boolean useRackAwareAssignment( + Set allMemberRacks, + Set allPartitionRacks, + Map> racksPerPartition + ) { + if (allMemberRacks.isEmpty() || Collections.disjoint(allMemberRacks, allPartitionRacks)) + return false; + else { + return !racksPerPartition.values().stream().allMatch(allPartitionRacks::equals); + } + } } diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java index 27ae8ce147b39..dc6e4b8cfc841 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java @@ -135,7 +135,7 @@ public class UniformHomogeneousAssignmentBuilder { .flatMap(Collection::stream) .collect(java.util.stream.Collectors.toSet()); - this.useRackStrategy = Utils.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition); + this.useRackStrategy = AssignorHelpers.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition); } /** @@ -224,7 +224,7 @@ private void maybeRevokePartitions() { // 1. We still have quota left to keep it. // 2-1. Don't use rack strategy, so we can keep it. // 2-2. Use rack strategy and the member rack matches the partition racks. - if (quota > 0 && (!useRackStrategy || Utils.isRackMatch( + if (quota > 0 && (!useRackStrategy || AssignorHelpers.isRackMatch( groupSpec.memberSubscription(memberId).rackId(), subscribedTopicDescriber.racksForPartition(topicId, partition)))) { quota--; @@ -298,7 +298,7 @@ private void assignRackAwarenessRemainingPartitions() { } String memberId = unfilledMember.memberId; - if (!Utils.isRackMatch(groupSpec.memberSubscription(memberId).rackId(), + if (!AssignorHelpers.isRackMatch(groupSpec.memberSubscription(memberId).rackId(), subscribedTopicDescriber.racksForPartition(tip.topicId(), tip.partitionId()))) { continue; } diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/Utils.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/Utils.java deleted file mode 100644 index cf47125528d7f..0000000000000 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/Utils.java +++ /dev/null @@ -1,42 +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.assignor; - -import org.apache.kafka.server.common.TopicIdPartition; - -import java.util.Collections; -import java.util.Map; -import java.util.Optional; -import java.util.Set; - -public class Utils { - public static boolean isRackMatch(Optional memberRackId, Set partitionRackIds) { - return memberRackId.isPresent() && partitionRackIds.contains(memberRackId.get()); - } - - public static boolean useRackAwareAssignment( - Set allMemberRacks, - Set allPartitionRacks, - Map> racksPerPartition - ) { - if (allMemberRacks.isEmpty() || Collections.disjoint(allMemberRacks, allPartitionRacks)) - return false; - else { - return !racksPerPartition.values().stream().allMatch(allPartitionRacks::equals); - } - } -} \ No newline at end of file diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/UtilsTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/AssignorHelpersTest.java similarity index 78% rename from group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/UtilsTest.java rename to group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/AssignorHelpersTest.java index 1f04074f22e33..df1762cd0461e 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/UtilsTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/AssignorHelpersTest.java @@ -28,27 +28,27 @@ import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertTrue; -public class UtilsTest { +public class AssignorHelpersTest { private static final Uuid TOPIC_ID = Uuid.randomUuid(); @Test void testIsRackMatchWithEmptyMemberRackId() { - assertFalse(Utils.isRackMatch(Optional.empty(), Set.of("rack1", "rack2"))); + assertFalse(AssignorHelpers.isRackMatch(Optional.empty(), Set.of("rack1", "rack2"))); } @Test void testIsRackMatchWithMatchingRack() { - assertTrue(Utils.isRackMatch(Optional.of("rack1"), Set.of("rack1", "rack2"))); + assertTrue(AssignorHelpers.isRackMatch(Optional.of("rack1"), Set.of("rack1", "rack2"))); } @Test void testIsRackMatchWithNonMatchingRack() { - assertFalse(Utils.isRackMatch(Optional.of("rack3"), Set.of("rack1", "rack2"))); + assertFalse(AssignorHelpers.isRackMatch(Optional.of("rack3"), Set.of("rack1", "rack2"))); } @Test void testIsRackMatchWithEmptyPartitionRacks() { - assertFalse(Utils.isRackMatch(Optional.of("rack1"), Set.of())); + assertFalse(AssignorHelpers.isRackMatch(Optional.of("rack1"), Set.of())); } @Test @@ -60,7 +60,7 @@ void testUseRackAwareAssignmentWithEmptyMemberRacks() { new TopicIdPartition(TOPIC_ID, 1), Set.of("rack2") ); - assertFalse(Utils.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition)); + assertFalse(AssignorHelpers.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition)); } @Test @@ -72,7 +72,7 @@ void testUseRackAwareAssignmentWithDisjointRacks() { new TopicIdPartition(TOPIC_ID, 1), Set.of("rack4") ); - assertFalse(Utils.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition)); + assertFalse(AssignorHelpers.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition)); } @Test @@ -84,7 +84,7 @@ void testUseRackAwareAssignmentWithAllPartitionsHavingSameRacks() { new TopicIdPartition(TOPIC_ID, 1), Set.of("rack1", "rack2") ); - assertFalse(Utils.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition)); + assertFalse(AssignorHelpers.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition)); } @Test @@ -96,6 +96,6 @@ void testUseRackAwareAssignmentWithDifferentRacksPerPartition() { new TopicIdPartition(TOPIC_ID, 1), Set.of("rack2") ); - assertTrue(Utils.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition)); + assertTrue(AssignorHelpers.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition)); } } From 92a61016c2d6bf70ed2657d832dcfe003a04a8ba Mon Sep 17 00:00:00 2001 From: PoAn Yang Date: Tue, 25 Nov 2025 22:01:53 +0800 Subject: [PATCH 3/6] use for-loop to collect member and partition racks Signed-off-by: PoAn Yang --- .../UniformHomogeneousAssignmentBuilder.java | 52 +++++++++---------- .../jmh/assignor/AssignorBenchmarkUtils.java | 5 ++ 2 files changed, 29 insertions(+), 28 deletions(-) diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java index dc6e4b8cfc841..35896532933c1 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java @@ -25,15 +25,12 @@ import org.apache.kafka.coordinator.group.modern.MemberAssignmentImpl; import org.apache.kafka.server.common.TopicIdPartition; -import java.util.ArrayList; -import java.util.Collection; import java.util.HashMap; import java.util.HashSet; import java.util.Iterator; import java.util.LinkedList; import java.util.List; import java.util.Map; -import java.util.Optional; import java.util.Set; /** @@ -96,6 +93,11 @@ public class UniformHomogeneousAssignmentBuilder { */ private final boolean useRackStrategy; + /** + * The mapping of topic partitions to their associated racks. + */ + private final Map> partitionRacks; + /** * The number of members to receive an extra partition beyond the minimum quota. * Example: If there are 11 partitions to be distributed among 3 members, @@ -112,30 +114,24 @@ public class UniformHomogeneousAssignmentBuilder { this.unassignedPartitions = new LinkedList<>(); this.targetAssignment = new HashMap<>(); + this.partitionRacks = new HashMap<>(); - Set allMemberRacks = groupSpec.memberIds().stream() - .map(memberId -> groupSpec.memberSubscription(memberId).rackId()) - .filter(Optional::isPresent) - .map(Optional::get) - .collect(java.util.stream.Collectors.toSet()); - Map> racksPerPartition = this.subscribedTopicIds.stream() - .flatMap(topicId -> { - int partitionCount = subscribedTopicDescriber.numPartitions(topicId); - List topicPartitions = new ArrayList<>(); - for (int partitionId = 0; partitionId < partitionCount; partitionId++) { - topicPartitions.add(new TopicIdPartition(topicId, partitionId)); - } - return topicPartitions.stream(); - }) - .collect(java.util.stream.Collectors.toMap( - tip -> tip, - tip -> subscribedTopicDescriber.racksForPartition(tip.topicId(), tip.partitionId()) - )); - Set allPartitionRacks = racksPerPartition.values().stream() - .flatMap(Collection::stream) - .collect(java.util.stream.Collectors.toSet()); - - this.useRackStrategy = AssignorHelpers.useRackAwareAssignment(allMemberRacks, allPartitionRacks, racksPerPartition); + Set allMemberRacks = new HashSet<>(); + for (String memberId : groupSpec.memberIds()) { + groupSpec.memberSubscription(memberId).rackId().ifPresent(allMemberRacks::add); + } + + Set allPartitionRacks = new HashSet<>(); + for (Uuid topicId : this.subscribedTopicIds) { + int partitionCount = subscribedTopicDescriber.numPartitions(topicId); + for (int partitionId = 0; partitionId < partitionCount; partitionId++) { + Set racks = subscribedTopicDescriber.racksForPartition(topicId, partitionId); + partitionRacks.put(new TopicIdPartition(topicId, partitionId), racks); + allPartitionRacks.addAll(racks); + } + } + + this.useRackStrategy = AssignorHelpers.useRackAwareAssignment(allMemberRacks, allPartitionRacks, partitionRacks); } /** @@ -226,7 +222,7 @@ private void maybeRevokePartitions() { // 2-2. Use rack strategy and the member rack matches the partition racks. if (quota > 0 && (!useRackStrategy || AssignorHelpers.isRackMatch( groupSpec.memberSubscription(memberId).rackId(), - subscribedTopicDescriber.racksForPartition(topicId, partition)))) { + partitionRacks.getOrDefault(new TopicIdPartition(topicId, partition), Set.of())))) { quota--; } else { if (newAssignment == null) { @@ -299,7 +295,7 @@ private void assignRackAwarenessRemainingPartitions() { String memberId = unfilledMember.memberId; if (!AssignorHelpers.isRackMatch(groupSpec.memberSubscription(memberId).rackId(), - subscribedTopicDescriber.racksForPartition(tip.topicId(), tip.partitionId()))) { + partitionRacks.getOrDefault(tip, Set.of()))) { continue; } diff --git a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/AssignorBenchmarkUtils.java b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/AssignorBenchmarkUtils.java index 555c92457f831..70108b7947bbd 100644 --- a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/AssignorBenchmarkUtils.java +++ b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/assignor/AssignorBenchmarkUtils.java @@ -18,6 +18,7 @@ import org.apache.kafka.common.Uuid; import org.apache.kafka.common.metadata.PartitionRecord; +import org.apache.kafka.common.metadata.RegisterBrokerRecord; import org.apache.kafka.common.metadata.TopicRecord; import org.apache.kafka.coordinator.common.runtime.CoordinatorMetadataImage; import org.apache.kafka.coordinator.common.runtime.KRaftCoordinatorMetadataImage; @@ -111,6 +112,10 @@ public static CoordinatorMetadataImage createMetadataImage( ); } + for (int i = 0; i < 4; i++) { + delta.replay(new RegisterBrokerRecord().setBrokerId(i).setRack("rack" + i)); + } + return new KRaftCoordinatorMetadataImage(delta.apply(MetadataProvenance.EMPTY)); } From 52aa53cf8d256245c57845ee34c340c454b098c8 Mon Sep 17 00:00:00 2001 From: PoAn Yang Date: Wed, 26 Nov 2025 20:23:25 +0800 Subject: [PATCH 4/6] feat: change back to use ArrayList Signed-off-by: PoAn Yang --- .../UniformHomogeneousAssignmentBuilder.java | 26 +++++++++---------- 1 file changed, 12 insertions(+), 14 deletions(-) diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java index 35896532933c1..26fec8a8346ce 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java @@ -25,10 +25,9 @@ import org.apache.kafka.coordinator.group.modern.MemberAssignmentImpl; import org.apache.kafka.server.common.TopicIdPartition; +import java.util.ArrayList; import java.util.HashMap; import java.util.HashSet; -import java.util.Iterator; -import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Set; @@ -69,7 +68,7 @@ public class UniformHomogeneousAssignmentBuilder { /** * The members that are below their quota. */ - private final LinkedList unfilledMembers; + private final List unfilledMembers; /** * The partitions that still need to be assigned. @@ -110,8 +109,8 @@ public class UniformHomogeneousAssignmentBuilder { this.subscribedTopicDescriber = subscribedTopicDescriber; this.subscribedTopicIds = new HashSet<>(groupSpec.memberSubscription(groupSpec.memberIds().iterator().next()) .subscribedTopicIds()); - this.unfilledMembers = new LinkedList<>(); - this.unassignedPartitions = new LinkedList<>(); + this.unfilledMembers = new ArrayList<>(); + this.unassignedPartitions = new ArrayList<>(); this.targetAssignment = new HashMap<>(); this.partitionRacks = new HashMap<>(); @@ -334,19 +333,18 @@ private void assignRemainingPartitions() { int unassignedPartitionIndex = 0; // If member is one of the last remainingMembersToGetAnExtraPartition, it must take the extra partition. - for (Iterator it = unfilledMembers.descendingIterator(); it.hasNext(); ) { - MemberWithRemainingQuota unfilledMember = it.next(); - if (remainingMembersToGetAnExtraPartition > 0) { - remainingMembersToGetAnExtraPartition--; - unfilledMember.incrementQuota(); - } - - if (unfilledMember.remainingQuota() == 0) { - it.remove(); + for (int i = unfilledMembers.size() - 1; i >= 0; i--) { + if (remainingMembersToGetAnExtraPartition == 0) { + break; } + unfilledMembers.get(i).incrementQuota(); + remainingMembersToGetAnExtraPartition--; } for (MemberWithRemainingQuota unfilledMember : unfilledMembers) { + if (unfilledMember.remainingQuota() == 0) { + continue; + } String memberId = unfilledMember.memberId; int remainingQuota = unfilledMember.remainingQuota; From 822a061c4f367725094be1ad7824949bb6aadad4 Mon Sep 17 00:00:00 2001 From: PoAn Yang Date: Thu, 27 Nov 2025 23:53:44 +0800 Subject: [PATCH 5/6] feat: assign rack awareness partitions with descending order Signed-off-by: PoAn Yang --- .../UniformHomogeneousAssignmentBuilder.java | 9 +++++---- .../OptimizedUniformAssignmentBuilderTest.java | 16 ++++++++-------- 2 files changed, 13 insertions(+), 12 deletions(-) diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java index 26fec8a8346ce..6d3d5e65bbe0c 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java @@ -74,7 +74,7 @@ public class UniformHomogeneousAssignmentBuilder { * The partitions that still need to be assigned. * Initially this contains all the subscribed topics' partitions. */ - private final List unassignedPartitions; + private List unassignedPartitions; /** * The target assignment. @@ -281,8 +281,9 @@ private void maybeRevokePartitions() { * Assign the unassigned partitions to the unfilled members if member and partition racks are matched. */ private void assignRackAwarenessRemainingPartitions() { - for (var unassignedPartitionsIter = unassignedPartitions.iterator(); unassignedPartitionsIter.hasNext(); ) { - TopicIdPartition tip = unassignedPartitionsIter.next(); + // Assign partitions to members with descending order. This avoids the cost of shifting elements. + for (int i = unassignedPartitions.size() - 1; i >= 0; i--) { + TopicIdPartition tip = unassignedPartitions.get(i); boolean isPartitionAssigned = false; for (var unfilledMembersIter = unfilledMembers.iterator(); unfilledMembersIter.hasNext(); ) { @@ -321,7 +322,7 @@ private void assignRackAwarenessRemainingPartitions() { } if (isPartitionAssigned) { - unassignedPartitionsIter.remove(); + unassignedPartitions.remove(i); } } } diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/OptimizedUniformAssignmentBuilderTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/OptimizedUniformAssignmentBuilderTest.java index 435556fda579c..8262b72bef2d6 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/OptimizedUniformAssignmentBuilderTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/OptimizedUniformAssignmentBuilderTest.java @@ -680,13 +680,13 @@ public void testFirstAssignmentThreeMembersOneTopicWithRacks() { Map>> expectedAssignment = new HashMap<>(); expectedAssignment.put(memberA, mkAssignment( - mkTopicAssignment(topic1Uuid, 0, 3) + mkTopicAssignment(topic1Uuid, 3, 4) )); expectedAssignment.put(memberB, mkAssignment( - mkTopicAssignment(topic1Uuid, 1, 4) + mkTopicAssignment(topic1Uuid, 1, 5) )); expectedAssignment.put(memberC, mkAssignment( - mkTopicAssignment(topic1Uuid, 2, 5) + mkTopicAssignment(topic1Uuid, 0, 2) )); assertAssignment(expectedAssignment, computedAssignment); @@ -847,15 +847,15 @@ public void testFirstAssignmentThreeMembersTwoTopicsWithOneMemberNoRack() { ); Map>> expectedAssignment = new HashMap<>(); expectedAssignment.put(memberA, mkAssignment( - mkTopicAssignment(topic1Uuid, 3, 4, 5) + mkTopicAssignment(topic1Uuid, 0, 3), + mkTopicAssignment(topic2Uuid, 0) )); expectedAssignment.put(memberB, mkAssignment( - mkTopicAssignment(topic1Uuid, 0), - mkTopicAssignment(topic2Uuid, 0, 1) + mkTopicAssignment(topic1Uuid, 1, 4, 5) )); expectedAssignment.put(memberC, mkAssignment( - mkTopicAssignment(topic1Uuid, 1, 2), - mkTopicAssignment(topic2Uuid, 2) + mkTopicAssignment(topic1Uuid, 2), + mkTopicAssignment(topic2Uuid, 1, 2) )); assertAssignment(expectedAssignment, computedAssignment); From 1590ab06455a9753a628916fdf8341172f5b4eea Mon Sep 17 00:00:00 2001 From: PoAn Yang Date: Mon, 8 Dec 2025 18:48:05 +0800 Subject: [PATCH 6/6] address comment Signed-off-by: PoAn Yang --- .../UniformHomogeneousAssignmentBuilder.java | 38 ++++++++++++------- .../group/assignor/CommonAssignorTests.java | 12 ++++-- ...iformHomogeneousAssignmentBuilderTest.java | 4 +- 3 files changed, 36 insertions(+), 18 deletions(-) diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java index 6d3d5e65bbe0c..849716bc06b4c 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilder.java @@ -26,6 +26,7 @@ import org.apache.kafka.server.common.TopicIdPartition; import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.List; @@ -65,6 +66,11 @@ public class UniformHomogeneousAssignmentBuilder { */ private final Set subscribedTopicIds; + /** + * The member Ids of the consumer group. + */ + private final List memberIds; + /** * The members that are below their quota. */ @@ -107,7 +113,9 @@ public class UniformHomogeneousAssignmentBuilder { UniformHomogeneousAssignmentBuilder(GroupSpec groupSpec, SubscribedTopicDescriber subscribedTopicDescriber) { this.groupSpec = groupSpec; this.subscribedTopicDescriber = subscribedTopicDescriber; - this.subscribedTopicIds = new HashSet<>(groupSpec.memberSubscription(groupSpec.memberIds().iterator().next()) + this.memberIds = new ArrayList<>(groupSpec.memberIds()); + Collections.sort(memberIds); + this.subscribedTopicIds = new HashSet<>(groupSpec.memberSubscription(this.memberIds.get(0)) .subscribedTopicIds()); this.unfilledMembers = new ArrayList<>(); this.unassignedPartitions = new ArrayList<>(); @@ -116,21 +124,25 @@ public class UniformHomogeneousAssignmentBuilder { this.partitionRacks = new HashMap<>(); Set allMemberRacks = new HashSet<>(); - for (String memberId : groupSpec.memberIds()) { + for (String memberId : this.memberIds) { groupSpec.memberSubscription(memberId).rackId().ifPresent(allMemberRacks::add); } - Set allPartitionRacks = new HashSet<>(); - for (Uuid topicId : this.subscribedTopicIds) { - int partitionCount = subscribedTopicDescriber.numPartitions(topicId); - for (int partitionId = 0; partitionId < partitionCount; partitionId++) { - Set racks = subscribedTopicDescriber.racksForPartition(topicId, partitionId); - partitionRacks.put(new TopicIdPartition(topicId, partitionId), racks); - allPartitionRacks.addAll(racks); + if (allMemberRacks.isEmpty()) { + this.useRackStrategy = false; + } else { + Set allPartitionRacks = new HashSet<>(); + for (Uuid topicId : this.subscribedTopicIds) { + int partitionCount = subscribedTopicDescriber.numPartitions(topicId); + for (int partitionId = 0; partitionId < partitionCount; partitionId++) { + Set racks = subscribedTopicDescriber.racksForPartition(topicId, partitionId); + partitionRacks.put(new TopicIdPartition(topicId, partitionId), racks); + allPartitionRacks.addAll(racks); + } } - } - this.useRackStrategy = AssignorHelpers.useRackAwareAssignment(allMemberRacks, allPartitionRacks, partitionRacks); + this.useRackStrategy = AssignorHelpers.useRackAwareAssignment(allMemberRacks, allPartitionRacks, partitionRacks); + } } /** @@ -161,7 +173,7 @@ public GroupAssignment build() throws PartitionAssignorException { // Compute the minimum required quota per member and the number of members // that should receive an extra partition. - int numberOfMembers = groupSpec.memberIds().size(); + int numberOfMembers = this.memberIds.size(); minimumMemberQuota = totalPartitionsCount / numberOfMembers; remainingMembersToGetAnExtraPartition = totalPartitionsCount % numberOfMembers; @@ -189,7 +201,7 @@ public GroupAssignment build() throws PartitionAssignorException { */ @SuppressWarnings({"CyclomaticComplexity", "NPathComplexity"}) private void maybeRevokePartitions() { - for (String memberId : groupSpec.memberIds()) { + for (String memberId : this.memberIds) { Map> oldAssignment = groupSpec.memberAssignment(memberId).partitions(); Map> newAssignment = null; diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/CommonAssignorTests.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/CommonAssignorTests.java index a7e4c4af675ba..7ea315cd0857d 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/CommonAssignorTests.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/CommonAssignorTests.java @@ -38,6 +38,7 @@ import static org.apache.kafka.coordinator.group.AssignmentTestUtil.assertAssignment; import static org.apache.kafka.coordinator.group.AssignmentTestUtil.invertedTargetAssignment; +import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertSame; public class CommonAssignorTests { @@ -123,9 +124,14 @@ public static void testAssignmentReuse(PartitionAssignor assignor, SubscriptionT ); for (String memberId : members.keySet()) { - // The assignment map from the assignor must be the same as the immutable assignment map - // that went in. - assertSame(membersWithAssignment.get(memberId).partitions(), secondAssignment.members().get(memberId).partitions()); + if (rackAware) { + // With rack awareness, the assignment maps may be mutable cause of revoking non-matched partitions. + assertEquals(membersWithAssignment.get(memberId).partitions(), secondAssignment.members().get(memberId).partitions()); + } else { + // The assignment map from the assignor must be the same as the immutable assignment map + // that went in. + assertSame(membersWithAssignment.get(memberId).partitions(), secondAssignment.members().get(memberId).partitions()); + } } } diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilderTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilderTest.java index 3235f3706f877..b04bf2a9da8b2 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilderTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/UniformHomogeneousAssignmentBuilderTest.java @@ -65,13 +65,13 @@ public class UniformHomogeneousAssignmentBuilderTest { private final String memberC = "C"; @ParameterizedTest - @ValueSource(booleans = {false}) + @ValueSource(booleans = {false, true}) public void testAssignmentReuse(boolean rackAware) { CommonAssignorTests.testAssignmentReuse(assignor, HOMOGENEOUS, rackAware); } @ParameterizedTest - @ValueSource(booleans = {false}) + @ValueSource(booleans = {false, true}) public void testReassignmentStickiness(boolean rackAware) { CommonAssignorTests.testReassignmentStickiness(assignor, HOMOGENEOUS, rackAware); }