84 lines
5.8 KiB
Diff
84 lines
5.8 KiB
Diff
# 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;
|
||
}
|