diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/KStream.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/KStream.java index 8df2131624eea..63e2439374dd9 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/KStream.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/KStream.java @@ -19,7 +19,7 @@ import org.apache.kafka.common.serialization.Deserializer; import org.apache.kafka.common.serialization.Serializer; -import org.apache.kafka.streaming.processor.ProcessorFactory; +import org.apache.kafka.streaming.processor.ProcessorDef; /** * KStream is an abstraction of a stream of key-value pairs. @@ -133,7 +133,7 @@ public interface KStream { /** * Processes all elements in this stream by applying a processor. * - * @param processorFactory the class of ProcessorFactory + * @param processorDef the class of ProcessorFactory */ - KStream process(ProcessorFactory processorFactory); + KStream process(ProcessorDef processorDef); } diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamBranch.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamBranch.java index e9d6030401c58..8fcfdfaca0e1e 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamBranch.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamBranch.java @@ -18,10 +18,10 @@ package org.apache.kafka.streaming.kstream.internals; import org.apache.kafka.streaming.processor.Processor; -import org.apache.kafka.streaming.processor.ProcessorFactory; +import org.apache.kafka.streaming.processor.ProcessorDef; import org.apache.kafka.streaming.kstream.Predicate; -class KStreamBranch implements ProcessorFactory { +class KStreamBranch implements ProcessorDef { private final Predicate[] predicates; diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamFilter.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamFilter.java index a85e0181a7c76..4f941feda66ed 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamFilter.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamFilter.java @@ -19,9 +19,9 @@ import org.apache.kafka.streaming.processor.Processor; import org.apache.kafka.streaming.kstream.Predicate; -import org.apache.kafka.streaming.processor.ProcessorFactory; +import org.apache.kafka.streaming.processor.ProcessorDef; -class KStreamFilter implements ProcessorFactory { +class KStreamFilter implements ProcessorDef { private final Predicate predicate; private final boolean filterOut; diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamFlatMap.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamFlatMap.java index bc834a8f66ae8..c5281f8a19dd7 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamFlatMap.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamFlatMap.java @@ -20,9 +20,9 @@ import org.apache.kafka.streaming.kstream.KeyValue; import org.apache.kafka.streaming.kstream.KeyValueFlatMap; import org.apache.kafka.streaming.processor.Processor; -import org.apache.kafka.streaming.processor.ProcessorFactory; +import org.apache.kafka.streaming.processor.ProcessorDef; -class KStreamFlatMap implements ProcessorFactory { +class KStreamFlatMap implements ProcessorDef { private final KeyValueFlatMap mapper; diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamFlatMapValues.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamFlatMapValues.java index 5e18f188ffa08..a04d1607d8be8 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamFlatMapValues.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamFlatMapValues.java @@ -19,9 +19,9 @@ import org.apache.kafka.streaming.processor.Processor; import org.apache.kafka.streaming.kstream.ValueMapper; -import org.apache.kafka.streaming.processor.ProcessorFactory; +import org.apache.kafka.streaming.processor.ProcessorDef; -class KStreamFlatMapValues implements ProcessorFactory { +class KStreamFlatMapValues implements ProcessorDef { private final ValueMapper> mapper; diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamImpl.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamImpl.java index e317338b134e6..567e47a0f464a 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamImpl.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamImpl.java @@ -20,7 +20,7 @@ import org.apache.kafka.common.serialization.Deserializer; import org.apache.kafka.common.serialization.Serializer; import org.apache.kafka.streaming.kstream.KeyValueFlatMap; -import org.apache.kafka.streaming.processor.ProcessorFactory; +import org.apache.kafka.streaming.processor.ProcessorDef; import org.apache.kafka.streaming.processor.TopologyBuilder; import org.apache.kafka.streaming.kstream.KStreamWindowed; import org.apache.kafka.streaming.kstream.KeyValueMapper; @@ -185,10 +185,10 @@ public void sendTo(String topic, Serializer keySerializer, Serializer valS @SuppressWarnings("unchecked") @Override - public KStream process(final ProcessorFactory processorFactory) { + public KStream process(final ProcessorDef processorDef) { String name = PROCESSOR_NAME + INDEX.getAndIncrement(); - topology.addProcessor(name, processorFactory, this.name); + topology.addProcessor(name, processorDef, this.name); return new KStreamImpl<>(topology, name); } diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamJoin.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamJoin.java index a8452d5ff0b64..3ec90048d5aad 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamJoin.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamJoin.java @@ -21,11 +21,11 @@ import org.apache.kafka.streaming.processor.ProcessorContext; import org.apache.kafka.streaming.kstream.ValueJoiner; import org.apache.kafka.streaming.kstream.Window; -import org.apache.kafka.streaming.processor.ProcessorFactory; +import org.apache.kafka.streaming.processor.ProcessorDef; import java.util.Iterator; -class KStreamJoin implements ProcessorFactory { +class KStreamJoin implements ProcessorDef { private static abstract class Finder { abstract Iterator find(K key, long timestamp); @@ -37,7 +37,7 @@ private static abstract class Finder { private final boolean prior; private Processor processorForOtherStream = null; - public final ProcessorFactory processorFactoryForOtherStream = new ProcessorFactory() { + public final ProcessorDef processorDefForOtherStream = new ProcessorDef() { @Override public Processor build() { return processorForOtherStream; diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamMap.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamMap.java index 83fb4c11c99f6..5399cc7a350b6 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamMap.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamMap.java @@ -20,9 +20,9 @@ import org.apache.kafka.streaming.processor.Processor; import org.apache.kafka.streaming.kstream.KeyValue; import org.apache.kafka.streaming.kstream.KeyValueMapper; -import org.apache.kafka.streaming.processor.ProcessorFactory; +import org.apache.kafka.streaming.processor.ProcessorDef; -class KStreamMap implements ProcessorFactory { +class KStreamMap implements ProcessorDef { private final KeyValueMapper mapper; diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamMapValues.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamMapValues.java index 667c04386c381..a1b589f236325 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamMapValues.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamMapValues.java @@ -19,9 +19,9 @@ import org.apache.kafka.streaming.processor.Processor; import org.apache.kafka.streaming.kstream.ValueMapper; -import org.apache.kafka.streaming.processor.ProcessorFactory; +import org.apache.kafka.streaming.processor.ProcessorDef; -class KStreamMapValues implements ProcessorFactory { +class KStreamMapValues implements ProcessorDef { private final ValueMapper mapper; diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamPassThrough.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamPassThrough.java index 98528f72c7c1b..6fa14f66b0f97 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamPassThrough.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamPassThrough.java @@ -18,9 +18,9 @@ package org.apache.kafka.streaming.kstream.internals; import org.apache.kafka.streaming.processor.Processor; -import org.apache.kafka.streaming.processor.ProcessorFactory; +import org.apache.kafka.streaming.processor.ProcessorDef; -class KStreamPassThrough implements ProcessorFactory { +class KStreamPassThrough implements ProcessorDef { @Override public Processor build() { diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamWindow.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamWindow.java index f357613d94ce4..018cdb309be93 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamWindow.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamWindow.java @@ -19,11 +19,11 @@ import org.apache.kafka.streaming.kstream.Window; import org.apache.kafka.streaming.processor.Processor; -import org.apache.kafka.streaming.processor.ProcessorFactory; +import org.apache.kafka.streaming.processor.ProcessorDef; import org.apache.kafka.streaming.processor.ProcessorContext; import org.apache.kafka.streaming.kstream.WindowDef; -public class KStreamWindow implements ProcessorFactory { +public class KStreamWindow implements ProcessorDef { private final WindowDef windowDef; diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamWindowedImpl.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamWindowedImpl.java index e269c3ad40bcc..62351a4e3607f 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamWindowedImpl.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamWindowedImpl.java @@ -37,7 +37,7 @@ private KStream join(KStreamWindowed other, boolean prior String joinName = JOIN_NAME + INDEX.getAndIncrement(); String joinOtherName = JOINOTHER_NAME + INDEX.getAndIncrement(); topology.addProcessor(joinName, join, this.name); - topology.addProcessor(joinOtherName, join.processorFactoryForOtherStream, ((KStreamImpl) other).name); + topology.addProcessor(joinOtherName, join.processorDefForOtherStream, ((KStreamImpl) other).name); return new KStreamImpl<>(topology, joinName); } diff --git a/stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorFactory.java b/stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorDef.java similarity index 96% rename from stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorFactory.java rename to stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorDef.java index 1d1c4816d8cce..642db734f5362 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorFactory.java +++ b/stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorDef.java @@ -17,7 +17,8 @@ package org.apache.kafka.streaming.processor; -public interface ProcessorFactory { +public interface ProcessorDef { Processor build(); + } 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 894e0e37f0244..5c39681e8a8f5 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 @@ -47,9 +47,9 @@ private interface NodeFactory { private class ProcessorNodeFactory implements NodeFactory { public final String[] parents; private final String name; - private final ProcessorFactory factory; + private final ProcessorDef factory; - public ProcessorNodeFactory(String name, String[] parents, ProcessorFactory factory) { + public ProcessorNodeFactory(String name, String[] parents, ProcessorDef factory) { this.name = name; this.parents = parents.clone(); this.factory = factory; @@ -134,7 +134,7 @@ public final void addSink(String name, String topic, Serializer keySerializer, S nodeFactories.add(new SinkNodeFactory(name, parentNames, topic, keySerializer, valSerializer)); } - public final void addProcessor(String name, ProcessorFactory factory, String... parentNames) { + public final void addProcessor(String name, ProcessorDef factory, String... parentNames) { if (nodeNames.contains(name)) throw new IllegalArgumentException("Processor " + name + " is already added."); @@ -186,7 +186,7 @@ public ProcessorTopology build() { processorMap.get(parent).addChild(node); } } else { - throw new IllegalStateException("unknown factory class: " + factory.getClass().getName()); + throw new IllegalStateException("unknown node factory class: " + factory.getClass().getName()); } } } catch (Exception e) {