Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;

/**
Expand Down Expand Up @@ -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<TopicPartition> sourceTopicPartitions();

/**
*
* @return the set of changelog topic partitions. This set will include both source and non-source
* topic partitions.
*/
Set<TopicPartition> changelogTopicPartitions();

/**
*
* @return the mapping of {@code TopicPartition} to set of rack ids that this partition resides on.
*/
Map<TopicPartition, Set<String>> partitionToRackIds();
Set<TaskTopicPartition> topicPartitions();
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -513,50 +515,45 @@ private ApplicationState buildApplicationState(final TopologyMetadata topologyMe
+ "tasks for source topics vs changelog topics.");
}

final Set<TopicPartition> sourceTopicPartitions = new HashSet<>();
final Set<TopicPartition> nonSourceChangelogTopicPartitions = new HashSet<>();
for (final Map.Entry<TaskId, Set<TopicPartition>> entry : sourcePartitionsForTask.entrySet()) {
final TaskId taskId = entry.getKey();
final Set<TopicPartition> taskSourcePartitions = entry.getValue();
final Set<TopicPartition> taskChangelogPartitions = changelogPartitionsForTask.get(taskId);
final Set<TopicPartition> taskNonSourceChangelogPartitions = new HashSet<>(taskChangelogPartitions);
taskNonSourceChangelogPartitions.removeAll(taskSourcePartitions);

sourceTopicPartitions.addAll(taskSourcePartitions);
nonSourceChangelogTopicPartitions.addAll(taskNonSourceChangelogPartitions);
}
final Set<TaskId> logicalTaskIds = unmodifiableSet(sourcePartitionsForTask.keySet());
final Set<DefaultTaskTopicPartition> allTopicPartitions = new HashSet<>();
final Map<TaskId, Set<TaskTopicPartition>> topicPartitionsForTask = new HashMap<>();
logicalTaskIds.forEach(taskId -> {
final Set<TaskTopicPartition> 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<TopicPartition, Set<String>> racksForSourcePartitions = RackUtils.getRacksForTopicPartition(
cluster, internalTopicManager, sourceTopicPartitions, false);
final Map<TopicPartition, Set<String>> 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<TaskId> logicalTaskIds = unmodifiableSet(sourcePartitionsForTask.keySet());
final Set<TaskInfo> logicalTasks = logicalTaskIds.stream().map(taskId -> {
final Set<String> stateStoreNames = topologyMetadata
.stateStoreNameToSourceTopicsForTopology(taskId.topologyName())
.keySet();
final Set<TopicPartition> sourcePartitions = sourcePartitionsForTask.get(taskId);
final Set<TopicPartition> changelogPartitions = changelogPartitionsForTask.get(taskId);
final Map<TopicPartition, Set<String>> 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<TaskTopicPartition> topicPartitions = topicPartitionsForTask.get(taskId);
return new DefaultTaskInfo(
taskId,
!stateStoreNames.isEmpty(),
racksForTaskPartition,
stateStoreNames,
sourcePartitions,
changelogPartitions
topicPartitions
);
}).collect(Collectors.toSet());

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<TopicPartition, Set<String>> partitionToRackIds;
private final Set<String> stateStoreNames;
private final Set<TopicPartition> sourceTopicPartitions;
private final Set<TopicPartition> changelogTopicPartitions;
private final Set<TaskTopicPartition> topicPartitions;

public DefaultTaskInfo(final TaskId id,
final boolean isStateful,
final Map<TopicPartition, Set<String>> partitionToRackIds,
final Set<String> stateStoreNames,
final Set<TopicPartition> sourceTopicPartitions,
final Set<TopicPartition> changelogTopicPartitions) {
final Set<TaskTopicPartition> 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
Expand All @@ -64,17 +56,7 @@ public Set<String> stateStoreNames() {
}

@Override
public Set<TopicPartition> sourceTopicPartitions() {
return sourceTopicPartitions;
}

@Override
public Set<TopicPartition> changelogTopicPartitions() {
return changelogTopicPartitions;
}

@Override
public Map<TopicPartition, Set<String>> partitionToRackIds() {
return partitionToRackIds;
public Set<TaskTopicPartition> topicPartitions() {
return topicPartitions;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -35,7 +36,9 @@ public class DefaultTaskTopicPartition implements TaskTopicPartition {
private final TopicPartition topicPartition;
private final boolean isSourceTopic;
private final boolean isChangelogTopic;
private final Optional<Set<String>> rackIds;

private Optional<Set<String>> rackIds;


public DefaultTaskTopicPartition(final TopicPartition topicPartition,
final boolean isSourceTopic,
Expand Down Expand Up @@ -66,4 +69,31 @@ public boolean isChangelog() {
public Optional<Set<String>> 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<String> rackIds) {
this.rackIds = Optional.of(rackIds);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -38,26 +39,37 @@ public final class RackUtils {

private RackUtils() { }

public static Map<TopicPartition, Set<String>> getRacksForTopicPartition(final Cluster cluster,
final InternalTopicManager internalTopicManager,
final Set<TopicPartition> topicPartitions,
final boolean isChangelog) {
final Set<String> 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<DefaultTaskTopicPartition> topicPartitions) {
// First we add all the changelog topics to the set of topics to describe.
final Set<String> topicsToDescribe = topicPartitions.stream()
.filter(DefaultTaskTopicPartition::isChangelog)
.map(topicPartition -> topicPartition.topicPartition().topic())
.collect(Collectors.toSet());

final Set<TopicPartition> 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<TopicPartition> nonChangelogTopics = topicPartitions.stream()
.filter(taskTopicPartition -> !taskTopicPartition.isChangelog())
.map(TaskTopicPartition::topicPartition)
.collect(Collectors.toSet());
final Map<TopicPartition, Set<String>> 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<String, List<TopicPartitionInfo>> freshTopicPartitionInfo =
describeTopics(internalTopicManager, topicsToDescribe);

// Finally we compute the list of topics that already have all rack information known.
final Set<TopicPartition> 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<TopicPartition, Set<String>> racksForTopicPartition = knownRacksForPartition(
cluster, topicsWithUpToDateMetadata);
freshTopicPartitionInfo.forEach((topic, partitionInfos) -> {
for (final TopicPartitionInfo partitionInfo : partitionInfos) {
final int partition = partitionInfo.partition();
Expand All @@ -75,7 +87,14 @@ public static Map<TopicPartition, Set<String>> getRacksForTopicPartition(final C
}
});

return racksForTopicPartition;
for (final DefaultTaskTopicPartition topicPartition : topicPartitions) {
if (!racksForTopicPartition.containsKey(topicPartition.topicPartition())) {
continue;
}

final Set<String> racks = racksForTopicPartition.get(topicPartition.topicPartition());
topicPartition.annotateWithRackIds(racks);
}
}

public static Set<String> topicsWithMissingMetadata(final Cluster cluster, final Set<TopicPartition> topicPartitions) {
Expand Down