java-topology/defects/kafka/patch/kafka-0001-0002-stickassignor-hashset.patch

84 lines
5.8 KiB
Diff
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# UNDF: UNDF-2026-000000131
--- 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
@@ -942,6 +942,9 @@ class AbstractStickyAssignor {
// a mapping of all topics to all consumers that can be assigned to them
private final Map<String, List<String>> topic2AllPotentialConsumers;
// a mapping of all consumers to all potential topics that can be assigned to them
- private final Map<String, List<String>> consumer2AllPotentialTopics;
+ // kafka-0002 fix: use Set<String> so contains() is O(1) instead of O(T).
+ // Previously Map<String, List<String>>; changed to Map<String, Set<String>> at
+ // construction time (line ~977) so maybeAssignPartition() gets O(1) lookup.
+ private final Map<String, Set<String>> consumer2AllPotentialTopics;
// a mapping of partition to current consumer
private final Map<TopicPartition, String> currentPartitionConsumer;
@@ -969,7 +972,7 @@ class AbstractStickyAssignor {
topic2AllPotentialConsumers = new HashMap<>(partitionsPerTopic.size());
- consumer2AllPotentialTopics = new HashMap<>(subscriptions.size());
+ consumer2AllPotentialTopics = new HashMap<>(subscriptions.size()); // values are now HashSet
// initialize topic2AllPotentialConsumers and consumer2AllPotentialTopics
partitionsPerTopic.keySet().forEach(
topicName -> topic2AllPotentialConsumers.put(topicName, new ArrayList<>()));
subscriptions.forEach((consumerId, subscription) -> {
- List<String> subscribedTopics = new ArrayList<>(subscription.topics().size());
+ // kafka-0002 fix: HashSet for O(1) contains() in maybeAssignPartition()
+ Set<String> subscribedTopics = new HashSet<>(subscription.topics().size() * 2);
consumer2AllPotentialTopics.put(consumerId, subscribedTopics);
@@ -1191,12 +1197,14 @@ class AbstractStickyAssignor {
for (String consumer: sortedCurrentSubscriptions) {
List<TopicPartition> consumerPartitions = currentAssignment.get(consumer);
int consumerPartitionCount = consumerPartitions.size();
// skip if this consumer already has all the topic partitions it can get
- List<String> allSubscribedTopics = consumer2AllPotentialTopics.get(consumer);
+ Set<String> allSubscribedTopics = consumer2AllPotentialTopics.get(consumer);
int maxAssignmentSize = getMaxAssignmentSize(allSubscribedTopics);
if (consumerPartitionCount == maxAssignmentSize)
continue;
+ // kafka-0001 fix: snapshot to HashSet<TopicPartition> so the inner
+ // contains() call is O(1) instead of O(P). Without this, the triple-
+ // nested loop (C × T × P) with an O(P) contains() gives O(C × T × P²).
+ Set<TopicPartition> consumerPartitionSet = new HashSet<>(consumerPartitions);
// otherwise make sure it cannot get any more
for (String topic: allSubscribedTopics) {
int partitionCount = partitionsPerTopic.get(topic).size();
for (int i = 0; i < partitionCount; i++) {
TopicPartition topicPartition = new TopicPartition(topic, i);
- if (!currentAssignment.get(consumer).contains(topicPartition)) {
+ if (!consumerPartitionSet.contains(topicPartition)) {
String otherConsumer = allPartitions.get(topicPartition);
int otherConsumerPartitionCount = currentAssignment.get(otherConsumer).size();
if (consumerPartitionCount + 1 < otherConsumerPartitionCount) {
@@ -1265,7 +1273,8 @@ class AbstractStickyAssignor {
private boolean maybeAssignPartition(TopicPartition partition, RackInfo rackInfo) {
for (String consumer: sortedCurrentSubscriptions) {
- if (consumer2AllPotentialTopics.get(consumer).contains(partition.topic()) && (rackInfo == null || !rackInfo.racksMismatch(consumer, partition))) {
+ // kafka-0002: consumer2AllPotentialTopics values are now HashSet<String> → O(1)
+ if (consumer2AllPotentialTopics.get(consumer).contains(partition.topic()) && (rackInfo == null || !rackInfo.racksMismatch(consumer, partition))) {
sortedCurrentSubscriptions.remove(consumer);
currentAssignment.get(consumer).add(partition);
currentPartitionConsumer.put(partition, consumer);
@@ -1303,7 +1312,7 @@ class AbstractStickyAssignor {
private boolean canConsumerParticipateInReassignment(String consumer) {
List<TopicPartition> currentPartitions = currentAssignment.get(consumer);
int currentAssignmentSize = currentPartitions.size();
- List<String> allSubscribedTopics = consumer2AllPotentialTopics.get(consumer);
+ Set<String> allSubscribedTopics = consumer2AllPotentialTopics.get(consumer);
int maxAssignmentSize = getMaxAssignmentSize(allSubscribedTopics);
@@ -1228,7 +1237,7 @@ class AbstractStickyAssignor {
- private int getMaxAssignmentSize(List<String> allSubscribedTopics) {
+ private int getMaxAssignmentSize(Set<String> allSubscribedTopics) {
int maxAssignmentSize;
if (allSubscribedTopics.size() == partitionsPerTopic.size()) {
maxAssignmentSize = totalPartitionsCount;
} else {
maxAssignmentSize = allSubscribedTopics.stream().map(partitionsPerTopic::get).map(List::size).reduce(0, Integer::sum);
}
return maxAssignmentSize;
}