java-topology/defects/kafka/patch/kafka-0010-roundrobin-assignor-topics-list-contains.patch

26 lines
1.7 KiB
Diff

# UNDF: UNDF-2026-000000686
diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/RoundRobinAssignor.java b/clients/src/main/java/org/apache/kafka/clients/consumer/RoundRobinAssignor.java
--- a/clients/src/main/java/org/apache/kafka/clients/consumer/RoundRobinAssignor.java
+++ b/clients/src/main/java/org/apache/kafka/clients/consumer/RoundRobinAssignor.java
@@ -103,6 +103,7 @@ public class RoundRobinAssignor extends AbstractPartitionAssignor {
List<MemberInfo> memberInfoList = new ArrayList<>();
for (Map.Entry<String, Subscription> memberSubscription : subscriptions.entrySet()) {
assignment.put(memberSubscription.getKey(), new ArrayList<>());
+ // Pre-convert topic list to Set so per-partition contains() is O(1) not O(T)
memberInfoList.add(new MemberInfo(memberSubscription.getKey(),
memberSubscription.getValue().groupInstanceId()));
}
@@ -112,7 +113,10 @@ public class RoundRobinAssignor extends AbstractPartitionAssignor {
CircularIterator<MemberInfo> assigner = new CircularIterator<>(Utils.sorted(memberInfoList));
+ Map<String, Set<String>> memberTopics = new HashMap<>();
+ subscriptions.forEach((memberId, sub) -> memberTopics.put(memberId, new HashSet<>(sub.topics())));
+
for (TopicPartition partition : allPartitionsSorted(partitionsPerTopic, subscriptions)) {
final String topic = partition.topic();
- while (!subscriptions.get(assigner.peek().memberId).topics().contains(topic))
+ while (!memberTopics.get(assigner.peek().memberId).contains(topic))
assigner.next();
assignment.get(assigner.next().memberId).add(partition);
}