diff --git a/stream/src/main/java/org/apache/kafka/streaming/KafkaStreaming.java b/stream/src/main/java/org/apache/kafka/streaming/KafkaStreaming.java index be27ae83acc7d..c1454aa916585 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/KafkaStreaming.java +++ b/stream/src/main/java/org/apache/kafka/streaming/KafkaStreaming.java @@ -22,8 +22,6 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.io.File; - /** * Kafka Streaming allows for performing continuous computation on input coming from one or more input topics and * sends output to zero or more output topics. @@ -118,5 +116,5 @@ public synchronized void close() { state = STOPPED; log.info("Stopped Kafka Stream process"); - } + } } 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 fee1f479390af..18ad1871dccea 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 @@ -31,47 +31,23 @@ private static abstract class Finder { abstract Iterator find(K key, long timestamp); } - private final String windowName1; - private final String windowName2; + private final String windowName; private final ValueJoiner joiner; - private Processor processorForOtherStream = null; - public final ProcessorDef processorDefForOtherStream = new ProcessorDef() { - @Override - public Processor instance() { - return processorForOtherStream; - } - }; - - KStreamJoin(String windowName1, String windowName2, ValueJoiner joiner) { - this.windowName1 = windowName1; - this.windowName2 = windowName2; + KStreamJoin(String windowName, ValueJoiner joiner) { + this.windowName = windowName; this.joiner = joiner; } @Override public Processor instance() { - // create a processor instance for the other stream - processorForOtherStream = new KStreamJoinProcessor(windowName1) { - @Override - protected void doJoin(K key, V2 value2, V1 value1) { - context.forward(key, joiner.apply(value1, value2)); - } - }; - - // create a processor instance for the primary stream - return new KStreamJoinProcessor(windowName2) { - @Override - protected void doJoin(K key, V1 value1, V2 value2) { - context.forward(key, joiner.apply(value1, value2)); - } - }; + return new KStreamJoinProcessor(windowName); } - private abstract class KStreamJoinProcessor extends KStreamProcessor { + private class KStreamJoinProcessor extends KStreamProcessor { private final String windowName; - protected Finder finder; + protected Finder finder; public KStreamJoinProcessor(String windowName) { this.windowName = windowName; @@ -86,27 +62,34 @@ public void init(ProcessorContext context) { if (!context.joinable()) throw new IllegalStateException("Streams are not joinable."); - final Window window = (Window) context.getStateStore(windowName); + final Window window = (Window) context.getStateStore(windowName); - this.finder = new Finder() { - Iterator find(K key, long timestamp) { + this.finder = new Finder() { + Iterator find(K key, long timestamp) { return window.find(key, timestamp); } }; } @Override - public void process(K key, T1 value) { + public void process(K key, V1 value) { long timestamp = context.timestamp(); - Iterator iter = finder.find(key, timestamp); + Iterator iter = finder.find(key, timestamp); if (iter != null) { while (iter.hasNext()) { - doJoin(key, value, iter.next()); + context.forward(key, joiner.apply(value, iter.next())); } } } + } - abstract protected void doJoin(K key, T1 value1, T2 value2); + public static ValueJoiner reserveJoiner(final ValueJoiner joiner) { + return new ValueJoiner() { + @Override + public R apply(T2 value2, T1 value1) { + return joiner.apply(value1, value2); + } + }; } } 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 7bbf45ee5e2c1..de2173b5f6e1f 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,15 +37,16 @@ public KStream join(KStreamWindowed other, ValueJoiner) other).windowDef.name(); + KStreamJoin joinThis = new KStreamJoin<>(otherWindowName, valueJoiner); + KStreamJoin joinOther = new KStreamJoin<>(thisWindowName, KStreamJoin.reserveJoiner(valueJoiner)); + KStreamPassThrough joinMerge = new KStreamPassThrough<>(); + String joinThisName = JOINTHIS_NAME + INDEX.getAndIncrement(); String joinOtherName = JOINOTHER_NAME + INDEX.getAndIncrement(); String joinMergeName = JOINMERGE_NAME + INDEX.getAndIncrement(); - KStreamJoin join = new KStreamJoin<>(thisWindowName, otherWindowName, valueJoiner); - KStreamPassThrough joinMerge = new KStreamPassThrough<>(); - - topology.addProcessor(joinThisName, join, this.name); - topology.addProcessor(joinOtherName, join.processorDefForOtherStream, ((KStreamImpl) other).name); + topology.addProcessor(joinThisName, joinThis, this.name); + topology.addProcessor(joinOtherName, joinOther, ((KStreamImpl) other).name); topology.addProcessor(joinMergeName, joinMerge, joinThisName, joinOtherName); return new KStreamImpl<>(topology, joinMergeName);