diff --git a/stream/src/main/java/org/apache/kafka/streaming/processor/TopologyBuilder.java b/stream/src/main/java/org/apache/kafka/streaming/processor/TopologyBuilder.java index 591b212963b28..6ee522a0dc0f3 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/processor/TopologyBuilder.java +++ b/stream/src/main/java/org/apache/kafka/streaming/processor/TopologyBuilder.java @@ -27,95 +27,130 @@ import java.util.ArrayList; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Set; public class TopologyBuilder { - private Map processorClasses = new HashMap<>(); - private Map sourceClasses = new HashMap<>(); - private Map sinkClasses = new HashMap<>(); - private Map topicsToSourceNames = new HashMap<>(); - private Map topicsToSinkNames = new HashMap<>(); + // list of node factories in a topological order + private ArrayList nodeFactories = new ArrayList<>(); - private Map> parents = new HashMap<>(); - private Map> children = new HashMap<>(); + private Set nodeNames = new HashSet<>(); + private Set sourceTopicNames = new HashSet<>(); - private class ProcessorClazz { - public ProcessorFactory factory; + private interface NodeFactory { + ProcessorNode build(); + } + + private class ProcessorNodeFactory implements NodeFactory { + public final String[] parents; + private final String name; + private final ProcessorFactory factory; - public ProcessorClazz(ProcessorFactory factory) { + public ProcessorNodeFactory(String name, String[] parents, ProcessorFactory factory) { + this.name = name; + this.parents = parents.clone(); this.factory = factory; } + + public ProcessorNode build() { + Processor processor = factory.build(); + return new ProcessorNode(name, processor); + } } - private class SourceClazz { - public Deserializer keyDeserializer; - public Deserializer valDeserializer; + private class SourceNodeFactory implements NodeFactory { + public final String[] topics; + private final String name; + private Deserializer keyDeserializer; + private Deserializer valDeserializer; - private SourceClazz(Deserializer keyDeserializer, Deserializer valDeserializer) { + private SourceNodeFactory(String name, String[] topics, Deserializer keyDeserializer, Deserializer valDeserializer) { + this.name = name; + this.topics = topics.clone(); this.keyDeserializer = keyDeserializer; this.valDeserializer = valDeserializer; } - } - private class SinkClazz { - public Serializer keySerializer; - public Serializer valSerializer; + public ProcessorNode build() { + return new SourceNode(name, keyDeserializer, valDeserializer); + } + } - private SinkClazz(Serializer keySerializer, Serializer valSerializer) { + private class SinkNodeFactory implements NodeFactory { + public final String[] parents; + public final String topic; + private final String name; + private Serializer keySerializer; + private Serializer valSerializer; + + private SinkNodeFactory(String name, String[] parents, String topic, Serializer keySerializer, Serializer valSerializer) { + this.name = name; + this.parents = parents.clone(); + this.topic = topic; this.keySerializer = keySerializer; this.valSerializer = valSerializer; } + public ProcessorNode build() { + return new SinkNode(name, topic, keySerializer, valSerializer); + } } public TopologyBuilder() {} - @SuppressWarnings("unchecked") public final void addSource(String name, Deserializer keyDeserializer, Deserializer valDeserializer, String... topics) { + if (nodeNames.contains(name)) + throw new IllegalArgumentException("Processor " + name + " is already added."); + for (String topic : topics) { - if (topicsToSourceNames.containsKey(topic)) + if (sourceTopicNames.contains(topic)) throw new IllegalArgumentException("Topic " + topic + " has already been registered by another processor."); - topicsToSourceNames.put(topic, name); + sourceTopicNames.add(topic); } - sourceClasses.put(name, new SourceClazz(keyDeserializer, valDeserializer)); + nodeNames.add(name); + nodeFactories.add(new SourceNodeFactory(name, topics, keyDeserializer, valDeserializer)); } - public final void addSink(String name, Serializer keySerializer, Serializer valSerializer, String... topics) { - for (String topic : topics) { - if (topicsToSinkNames.containsKey(topic)) - throw new IllegalArgumentException("Topic " + topic + " has already been registered by another processor."); + public final void addSink(String name, String topic, Serializer keySerializer, Serializer valSerializer, String... parentNames) { + if (nodeNames.contains(name)) + throw new IllegalArgumentException("Processor " + name + " is already added."); - topicsToSinkNames.put(topic, name); + if (parentNames != null) { + for (String parent : parentNames) { + if (parent.equals(name)) { + throw new IllegalArgumentException("Processor " + name + " cannot be a parent of itself"); + } + if (!nodeNames.contains(parent)) { + throw new IllegalArgumentException("Parent processor " + parent + " is not added yet."); + } + } } - sinkClasses.put(name, new SinkClazz(keySerializer, valSerializer)); + nodeNames.add(name); + nodeFactories.add(new SinkNodeFactory(name, parentNames, topic, keySerializer, valSerializer)); } public final void addProcessor(String name, ProcessorFactory factory, String... parentNames) { - if (processorClasses.containsKey(name)) + if (nodeNames.contains(name)) throw new IllegalArgumentException("Processor " + name + " is already added."); - processorClasses.put(name, new ProcessorClazz(factory)); - if (parentNames != null) { for (String parent : parentNames) { - if (!processorClasses.containsKey(parent)) + if (parent.equals(name)) { + throw new IllegalArgumentException("Processor " + name + " cannot be a parent of itself"); + } + if (!nodeNames.contains(parent)) { throw new IllegalArgumentException("Parent processor " + parent + " is not added yet."); - - // add to parent list - if (!parents.containsKey(name)) - parents.put(name, new ArrayList<>()); - parents.get(name).add(parent); - - // add to children list - if (!children.containsKey(parent)) - children.put(parent, new ArrayList<>()); - children.get(parent).add(name); + } } } + + nodeNames.add(name); + nodeFactories.add(new ProcessorNodeFactory(name, parentNames, factory)); } /** @@ -123,59 +158,41 @@ public final void addProcessor(String name, ProcessorFactory factory, String... */ @SuppressWarnings("unchecked") public ProcessorTopology build() { + List processorNodes = new ArrayList<>(nodeFactories.size()); Map processorMap = new HashMap<>(); Map topicSourceMap = new HashMap<>(); Map topicSinkMap = new HashMap<>(); - // create sources - for (String name : sourceClasses.keySet()) { - Deserializer keyDeserializer = sourceClasses.get(name).keyDeserializer; - Deserializer valDeserializer = sourceClasses.get(name).valDeserializer; - SourceNode node = new SourceNode(name, keyDeserializer, valDeserializer); - processorMap.put(name, node); - } - - // create sinks - for (String name : sinkClasses.keySet()) { - Serializer keySerializer = sinkClasses.get(name).keySerializer; - Serializer valSerializer = sinkClasses.get(name).valSerializer; - SinkNode node = new SinkNode(name, keySerializer, valSerializer); - processorMap.put(name, node); - } - - // create processors try { - for (String name : processorClasses.keySet()) { - ProcessorFactory processorFactory = processorClasses.get(name).factory; - Processor processor = processorFactory.build(); - ProcessorNode node = new ProcessorNode(name, processor); - processorMap.put(name, node); + // create processor nodes in a topological order ("nodeFactories" is already topologically sorted) + for (NodeFactory factory : nodeFactories) { + ProcessorNode node = factory.build(); + processorNodes.add(node); + processorMap.put(node.name(), node); + + if (factory instanceof ProcessorNodeFactory) { + for (String parent : ((ProcessorNodeFactory) factory).parents) { + processorMap.get(parent).chain(node); + } + } else if (factory instanceof SourceNodeFactory) { + for (String topic : ((SourceNodeFactory)factory).topics) { + topicSourceMap.put(topic, (SourceNode) node); + } + } else if (factory instanceof SinkNodeFactory) { + String topic = ((SinkNodeFactory) factory).topic; + topicSinkMap.put(topic, (SinkNode) node); + + for (String parent : ((SinkNodeFactory) factory).parents) { + processorMap.get(parent).chain(node); + } + } else { + throw new IllegalStateException("unknown factory class: " + factory.getClass().getName()); + } } } catch (Exception e) { - throw new KafkaException("Processor(String) constructor failed: this should not happen."); - } - - // construct topics to sources map - for (String topic : topicsToSourceNames.keySet()) { - SourceNode node = (SourceNode) processorMap.get(topicsToSourceNames.get(topic)); - topicSourceMap.put(topic, node); - } - - // construct topics to sinks map - for (String topic : topicsToSinkNames.keySet()) { - SinkNode node = (SinkNode) processorMap.get(topicsToSourceNames.get(topic)); - topicSinkMap.put(topic, node); - node.addTopic(topic); - } - - // chain children to parents to build the DAG - for (ProcessorNode node : processorMap.values()) { - for (String child : children.get(node.name())) { - ProcessorNode childNode = processorMap.get(child); - node.chain(childNode); - } + throw new KafkaException("ProcessorNode construction failed: this should not happen."); } - return new ProcessorTopology(processorMap, topicSourceMap, topicSinkMap); + return new ProcessorTopology(processorNodes, topicSourceMap, topicSinkMap); } } diff --git a/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorNode.java b/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorNode.java index 167863172149e..480c577754fdb 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorNode.java +++ b/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorNode.java @@ -31,8 +31,6 @@ public class ProcessorNode { private final String name; private final Processor processor; - public boolean initialized; - public ProcessorNode(String name) { this(name, null); } @@ -42,8 +40,6 @@ public ProcessorNode(String name, Processor processor) { this.processor = processor; this.parents = new ArrayList<>(); this.children = new ArrayList<>(); - - this.initialized = false; } public String name() { diff --git a/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorTopology.java b/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorTopology.java index 6ae7c3175096d..bf4ddd796ece8 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorTopology.java +++ b/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorTopology.java @@ -27,88 +27,60 @@ public class ProcessorTopology { - private Map processors = new HashMap<>(); - private Map sourceTopics = new HashMap<>(); - private Map sinkTopics = new HashMap<>(); + private List processors; + private Map sourceByTopics; + private Map sinkByTopics; - public ProcessorTopology(Map processors, - Map sourceTopics, - Map sinkTopics) { + public ProcessorTopology(List processors, + Map sourceByTopics, + Map sinkByTopics) { this.processors = processors; - this.sourceTopics = sourceTopics; - this.sinkTopics = sinkTopics; + this.sourceByTopics = sourceByTopics; + this.sinkByTopics = sinkByTopics; } public Set sourceTopics() { - return sourceTopics.keySet(); + return sourceByTopics.keySet(); } public Set sinkTopics() { - return sinkTopics.keySet(); + return sinkByTopics.keySet(); } public SourceNode source(String topic) { - return sourceTopics.get(topic); + return sourceByTopics.get(topic); } public SinkNode sink(String topic) { - return sinkTopics.get(topic); + return sinkByTopics.get(topic); } public Collection sources() { - return sourceTopics.values(); + return sourceByTopics.values(); } public Collection sinks() { - return sinkTopics.values(); + return sinkByTopics.values(); } /** * Initialize the processors following the DAG reverse ordering * such that parents are always initialized before children */ - @SuppressWarnings("unchecked") public void init(ProcessorContext context) { - // initialize sources - for (String topic : sourceTopics.keySet()) { - SourceNode source = sourceTopics.get(topic); - - init(source, context); - } - } - - /** - * Initialize the current processor node by first initializing - * its parent nodes first, then the processor itself - */ - @SuppressWarnings("unchecked") - private void init(ProcessorNode node, ProcessorContext context) { - for (ProcessorNode parentNode : (List) node.parents()) { - if (!parentNode.initialized) { - init(parentNode, context); - } - } - - node.init(context); - node.initialized = true; - - // try to initialize its children - for (ProcessorNode childNode : (List) node.children()) { - if (!childNode.initialized) { - init(childNode, context); - } + for (ProcessorNode node : processors) { + node.init(context); } } public final void close() { // close the processors - // TODO: do we need to follow the DAG ordering? - for (ProcessorNode processorNode : processors.values()) { - processorNode.close(); + for (ProcessorNode node : processors) { + node.close(); } processors.clear(); - sourceTopics.clear(); - sinkTopics.clear(); + sourceByTopics.clear(); + sinkByTopics.clear(); } } diff --git a/stream/src/main/java/org/apache/kafka/streaming/processor/internals/SinkNode.java b/stream/src/main/java/org/apache/kafka/streaming/processor/internals/SinkNode.java index f3cbadf2deda4..0daa054657130 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/processor/internals/SinkNode.java +++ b/stream/src/main/java/org/apache/kafka/streaming/processor/internals/SinkNode.java @@ -26,24 +26,20 @@ public class SinkNode extends ProcessorNode { + private final String topic; private final Serializer keySerializer; private final Serializer valSerializer; - private final List topics; private ProcessorContext context; - public SinkNode(String name, Serializer keySerializer, Serializer valSerializer) { + public SinkNode(String name, String topic, Serializer keySerializer, Serializer valSerializer) { super(name); - this.topics = new ArrayList<>(); + this.topic = topic; this.keySerializer = keySerializer; this.valSerializer = valSerializer; } - public void addTopic(String topic) { - this.topics.add(topic); - } - @Override public void init(ProcessorContext context) { this.context = context; @@ -53,9 +49,7 @@ public void init(ProcessorContext context) { public void process(K key, V value) { // send to all the registered topics RecordCollector collector = ((ProcessorContextImpl)context).recordCollector(); - for (String topic : topics) { - collector.send(new ProducerRecord<>(topic, key, value), keySerializer, valSerializer); - } + collector.send(new ProducerRecord<>(topic, key, value), keySerializer, valSerializer); } @Override