java-topology/defects/kafka/patch/kafka-0011-sticky-assignor-topics-list-contains.patch

31 lines
2.5 KiB
Diff

# UNDF: UNDF-2026-000000687
diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignor.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignor.java
--- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignor.java
+++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignor.java
@@ -577,6 +577,7 @@ public abstract class AbstractStickyAssignor extends AbstractPartitionAssignor {
private final Map<String, List<TopicPartition>> currentAssignment;
private final Map<String, Subscription> subscriptions;
private final Map<String, List<String>> consumer2AllPotentialTopics;
+ private final Map<String, Set<String>> consumer2TopicsSet;
private final Map<TopicPartition, List<String>> partition2AllPotentialConsumers;
private final RackInfo rackInfo;
private final int minQuota;
@@ -615,6 +616,9 @@ public abstract class AbstractStickyAssignor extends AbstractPartitionAssignor {
Set<TopicPartition> partitionsWithMultiplePreviousOwners) {
this.currentAssignment = currentAssignment;
this.subscriptions = subscriptions;
+ // Build a Set-backed view so that topics().contains() in assignOwnedPartitions is O(1)
+ this.consumer2TopicsSet = new HashMap<>();
+ subscriptions.forEach((consumer, sub) -> consumer2TopicsSet.put(consumer, new HashSet<>(sub.topics())));
this.consumer2AllPotentialTopics = consumer2AllPotentialTopics;
this.partition2AllPotentialConsumers = partition2AllPotentialConsumers;
this.rackInfo = rackInfo;
@@ -1045,7 +1049,7 @@ public abstract class AbstractStickyAssignor extends AbstractPartitionAssignor {
if (!topic2AllPotentialConsumers.containsKey(partition.topic())) {
partitionIter.remove();
currentPartitionConsumer.remove(partition);
- } else if (!consumerSubscription.topics().contains(partition.topic()) || rackInfo.racksMismatch(consumer, partition)) {
+ } else if (!consumer2TopicsSet.get(consumer).contains(partition.topic()) || rackInfo.racksMismatch(consumer, partition)) {
partitionIter.remove();
revocationRequired = true;
} else {