diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/assignment/TaskInfo.java b/streams/src/main/java/org/apache/kafka/streams/processor/assignment/TaskInfo.java index 2e61c2a5b561a..e8fe8b014d667 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/assignment/TaskInfo.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/assignment/TaskInfo.java @@ -17,9 +17,7 @@ package org.apache.kafka.streams.processor.assignment; -import java.util.Map; import java.util.Set; -import org.apache.kafka.common.TopicPartition; import org.apache.kafka.streams.processor.TaskId; /** @@ -52,21 +50,7 @@ public interface TaskInfo { /** * - * @return the set of source topic partitions. This set will include both changelog and non-changelog - * topic partitions. + * @return the set of topic partitions in use for this task. */ - Set sourceTopicPartitions(); - - /** - * - * @return the set of changelog topic partitions. This set will include both source and non-source - * topic partitions. - */ - Set changelogTopicPartitions(); - - /** - * - * @return the mapping of {@code TopicPartition} to set of rack ids that this partition resides on. - */ - Map> partitionToRackIds(); + Set topicPartitions(); } \ No newline at end of file diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java index 82a10615293e4..d7c15f9fd0d69 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java @@ -41,6 +41,7 @@ import org.apache.kafka.streams.processor.assignment.TaskInfo; import org.apache.kafka.streams.processor.assignment.ProcessId; import org.apache.kafka.streams.processor.assignment.TaskAssignor.TaskAssignment; +import org.apache.kafka.streams.processor.assignment.TaskTopicPartition; import org.apache.kafka.streams.processor.internals.assignment.ApplicationStateImpl; import org.apache.kafka.streams.processor.internals.InternalTopologyBuilder.TopicsInfo; import org.apache.kafka.streams.processor.internals.TopologyMetadata.Subtopology; @@ -51,6 +52,7 @@ import org.apache.kafka.streams.processor.internals.assignment.AssignorError; import org.apache.kafka.streams.processor.internals.assignment.ClientState; import org.apache.kafka.streams.processor.internals.assignment.CopartitionedTopicsEnforcer; +import org.apache.kafka.streams.processor.internals.assignment.DefaultTaskTopicPartition; import org.apache.kafka.streams.processor.internals.assignment.FallbackPriorTaskAssignor; import org.apache.kafka.streams.processor.internals.assignment.RackAwareTaskAssignor; import org.apache.kafka.streams.processor.internals.assignment.RackUtils; @@ -513,50 +515,45 @@ private ApplicationState buildApplicationState(final TopologyMetadata topologyMe + "tasks for source topics vs changelog topics."); } - final Set sourceTopicPartitions = new HashSet<>(); - final Set nonSourceChangelogTopicPartitions = new HashSet<>(); - for (final Map.Entry> entry : sourcePartitionsForTask.entrySet()) { - final TaskId taskId = entry.getKey(); - final Set taskSourcePartitions = entry.getValue(); - final Set taskChangelogPartitions = changelogPartitionsForTask.get(taskId); - final Set taskNonSourceChangelogPartitions = new HashSet<>(taskChangelogPartitions); - taskNonSourceChangelogPartitions.removeAll(taskSourcePartitions); - - sourceTopicPartitions.addAll(taskSourcePartitions); - nonSourceChangelogTopicPartitions.addAll(taskNonSourceChangelogPartitions); - } + final Set logicalTaskIds = unmodifiableSet(sourcePartitionsForTask.keySet()); + final Set allTopicPartitions = new HashSet<>(); + final Map> topicPartitionsForTask = new HashMap<>(); + logicalTaskIds.forEach(taskId -> { + final Set topicPartitions = new HashSet<>(); + + for (final TopicPartition topicPartition : sourcePartitionsForTask.get(taskId)) { + final boolean isSource = true; + final boolean isChangelog = changelogPartitionsForTask.get(taskId).contains(topicPartition); + final DefaultTaskTopicPartition racklessTopicPartition = new DefaultTaskTopicPartition( + topicPartition, isSource, isChangelog, null); + allTopicPartitions.add(racklessTopicPartition); + topicPartitions.add(racklessTopicPartition); + } - final Map> racksForSourcePartitions = RackUtils.getRacksForTopicPartition( - cluster, internalTopicManager, sourceTopicPartitions, false); - final Map> racksForChangelogPartitions = RackUtils.getRacksForTopicPartition( - cluster, internalTopicManager, nonSourceChangelogTopicPartitions, true); + for (final TopicPartition topicPartition : changelogPartitionsForTask.get(taskId)) { + final boolean isSource = sourcePartitionsForTask.get(taskId).contains(topicPartition); + final boolean isChangelog = true; + final DefaultTaskTopicPartition racklessTopicPartition = new DefaultTaskTopicPartition( + topicPartition, isSource, isChangelog, null); + allTopicPartitions.add(racklessTopicPartition); + topicPartitions.add(racklessTopicPartition); + } + + topicPartitionsForTask.put(taskId, topicPartitions); + }); + + RackUtils.annotateTopicPartitionsWithRackInfo(cluster, internalTopicManager, allTopicPartitions); - final Set logicalTaskIds = unmodifiableSet(sourcePartitionsForTask.keySet()); final Set logicalTasks = logicalTaskIds.stream().map(taskId -> { final Set stateStoreNames = topologyMetadata .stateStoreNameToSourceTopicsForTopology(taskId.topologyName()) .keySet(); - final Set sourcePartitions = sourcePartitionsForTask.get(taskId); - final Set changelogPartitions = changelogPartitionsForTask.get(taskId); - final Map> racksForTaskPartition = new HashMap<>(); - sourcePartitions.forEach(topicPartition -> { - racksForTaskPartition.put(topicPartition, racksForSourcePartitions.get(topicPartition)); - }); - changelogPartitions.forEach(topicPartition -> { - if (racksForSourcePartitions.containsKey(topicPartition)) { - racksForTaskPartition.put(topicPartition, racksForSourcePartitions.get(topicPartition)); - } else { - racksForTaskPartition.put(topicPartition, racksForChangelogPartitions.get(topicPartition)); - } - }); - + final Set topicPartitions = topicPartitionsForTask.get(taskId); return new DefaultTaskInfo( taskId, !stateStoreNames.isEmpty(), - racksForTaskPartition, stateStoreNames, - sourcePartitions, - changelogPartitions + topicPartitions ); }).collect(Collectors.toSet()); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/DefaultTaskInfo.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/DefaultTaskInfo.java index c0212db862af2..14e53440f4bf8 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/DefaultTaskInfo.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/DefaultTaskInfo.java @@ -16,36 +16,28 @@ */ package org.apache.kafka.streams.processor.internals.assignment; -import static java.util.Collections.unmodifiableMap; import static java.util.Collections.unmodifiableSet; -import java.util.Map; import java.util.Set; -import org.apache.kafka.common.TopicPartition; import org.apache.kafka.streams.processor.TaskId; import org.apache.kafka.streams.processor.assignment.TaskInfo; +import org.apache.kafka.streams.processor.assignment.TaskTopicPartition; public class DefaultTaskInfo implements TaskInfo { private final TaskId id; private final boolean isStateful; - private final Map> partitionToRackIds; private final Set stateStoreNames; - private final Set sourceTopicPartitions; - private final Set changelogTopicPartitions; + private final Set topicPartitions; public DefaultTaskInfo(final TaskId id, final boolean isStateful, - final Map> partitionToRackIds, final Set stateStoreNames, - final Set sourceTopicPartitions, - final Set changelogTopicPartitions) { + final Set topicPartitions) { this.id = id; - this.partitionToRackIds = unmodifiableMap(partitionToRackIds); this.isStateful = isStateful; this.stateStoreNames = unmodifiableSet(stateStoreNames); - this.sourceTopicPartitions = unmodifiableSet(sourceTopicPartitions); - this.changelogTopicPartitions = unmodifiableSet(changelogTopicPartitions); + this.topicPartitions = unmodifiableSet(topicPartitions); } @Override @@ -64,17 +56,7 @@ public Set stateStoreNames() { } @Override - public Set sourceTopicPartitions() { - return sourceTopicPartitions; - } - - @Override - public Set changelogTopicPartitions() { - return changelogTopicPartitions; - } - - @Override - public Map> partitionToRackIds() { - return partitionToRackIds; + public Set topicPartitions() { + return topicPartitions; } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/DefaultTaskTopicPartition.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/DefaultTaskTopicPartition.java index 815aa1ff64c92..1c0a640d9c4c9 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/DefaultTaskTopicPartition.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/DefaultTaskTopicPartition.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.streams.processor.internals.assignment; +import java.util.Objects; import java.util.Optional; import java.util.Set; import org.apache.kafka.common.TopicPartition; @@ -35,7 +36,9 @@ public class DefaultTaskTopicPartition implements TaskTopicPartition { private final TopicPartition topicPartition; private final boolean isSourceTopic; private final boolean isChangelogTopic; - private final Optional> rackIds; + + private Optional> rackIds; + public DefaultTaskTopicPartition(final TopicPartition topicPartition, final boolean isSourceTopic, @@ -66,4 +69,31 @@ public boolean isChangelog() { public Optional> rackIds() { return rackIds; } + + @Override + public int hashCode() { + int result = topicPartition.hashCode(); + result = 31 * result + Objects.hashCode(isSourceTopic); + result = 31 * result + Objects.hashCode(isChangelogTopic); + return result; + } + + @Override + public boolean equals(final Object obj) { + if (this == obj) + return true; + if (obj == null) + return false; + if (getClass() != obj.getClass()) + return false; + final TaskTopicPartition other = (TaskTopicPartition) obj; + return topicPartition.equals(other.topicPartition()) && + isSourceTopic == other.isSource() && + isChangelogTopic == other.isChangelog() && + rackIds.equals(other.rackIds()); + } + + public void annotateWithRackIds(final Set rackIds) { + this.rackIds = Optional.of(rackIds); + } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackUtils.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackUtils.java index b3554c36b03b1..f9961c555fa6c 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackUtils.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/assignment/RackUtils.java @@ -28,6 +28,7 @@ import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.TopicPartitionInfo; +import org.apache.kafka.streams.processor.assignment.TaskTopicPartition; import org.apache.kafka.streams.processor.internals.InternalTopicManager; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -38,26 +39,37 @@ public final class RackUtils { private RackUtils() { } - public static Map> getRacksForTopicPartition(final Cluster cluster, - final InternalTopicManager internalTopicManager, - final Set topicPartitions, - final boolean isChangelog) { - final Set topicsToDescribe = new HashSet<>(); - if (isChangelog) { - topicsToDescribe.addAll(topicPartitions.stream().map(TopicPartition::topic).collect( - Collectors.toSet())); - } else { - topicsToDescribe.addAll(topicsWithMissingMetadata(cluster, topicPartitions)); - } + public static void annotateTopicPartitionsWithRackInfo(final Cluster cluster, + final InternalTopicManager internalTopicManager, + final Set topicPartitions) { + // First we add all the changelog topics to the set of topics to describe. + final Set topicsToDescribe = topicPartitions.stream() + .filter(DefaultTaskTopicPartition::isChangelog) + .map(topicPartition -> topicPartition.topicPartition().topic()) + .collect(Collectors.toSet()); - final Set topicsWithUpToDateMetadata = topicPartitions.stream() - .filter(partition -> !topicsToDescribe.contains(partition.topic())) + // Then we add the non changelog topics that we do not have full information about. + final Set nonChangelogTopics = topicPartitions.stream() + .filter(taskTopicPartition -> !taskTopicPartition.isChangelog()) + .map(TaskTopicPartition::topicPartition) .collect(Collectors.toSet()); - final Map> racksForTopicPartition = knownRacksForPartition( - cluster, topicsWithUpToDateMetadata); + topicsToDescribe.addAll(topicsWithMissingMetadata(cluster, nonChangelogTopics)); + // We can issue an RPC call to get up-to-date information about the topics that had rack + // information missing. final Map> freshTopicPartitionInfo = describeTopics(internalTopicManager, topicsToDescribe); + + // Finally we compute the list of topics that already have all rack information known. + final Set topicsWithUpToDateMetadata = topicPartitions.stream() + .map(TaskTopicPartition::topicPartition) + .filter(topicPartition -> !topicsToDescribe.contains(topicPartition.topic())) + .collect(Collectors.toSet()); + + // Lastly we compile the mapping of topic partition to rack ids by combining known data and + // information that we got from the earlier RPC call. + final Map> racksForTopicPartition = knownRacksForPartition( + cluster, topicsWithUpToDateMetadata); freshTopicPartitionInfo.forEach((topic, partitionInfos) -> { for (final TopicPartitionInfo partitionInfo : partitionInfos) { final int partition = partitionInfo.partition(); @@ -75,7 +87,14 @@ public static Map> getRacksForTopicPartition(final C } }); - return racksForTopicPartition; + for (final DefaultTaskTopicPartition topicPartition : topicPartitions) { + if (!racksForTopicPartition.containsKey(topicPartition.topicPartition())) { + continue; + } + + final Set racks = racksForTopicPartition.get(topicPartition.topicPartition()); + topicPartition.annotateWithRackIds(racks); + } } public static Set topicsWithMissingMetadata(final Cluster cluster, final Set topicPartitions) {