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 fc23e2d6d530d..8df2131624eea 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.Processor; +import org.apache.kafka.streaming.processor.ProcessorFactory; /** * KStream is an abstraction of a stream of key-value pairs. @@ -83,10 +83,10 @@ public interface KStream { /** * Creates a new windowed stream using a specified window instance. * - * @param window the instance of Window + * @param windowDef the instance of Window * @return KStream */ - KStreamWindowed with(Window window); + KStreamWindowed with(WindowDef windowDef); /** * Creates an array of streams from this stream. Each stream in the array coresponds to a predicate in @@ -133,7 +133,7 @@ public interface KStream { /** * Processes all elements in this stream by applying a processor. * - * @param processor the class of Processor + * @param processorFactory the class of ProcessorFactory */ - KStream process(Processor processor); + KStream process(ProcessorFactory processorFactory); } diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/SlidingWindow.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/SlidingWindow.java deleted file mode 100644 index ad50f436300e3..0000000000000 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/SlidingWindow.java +++ /dev/null @@ -1,251 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.kafka.streaming.kstream; - -import org.apache.kafka.clients.producer.ProducerRecord; -import org.apache.kafka.common.serialization.ByteArraySerializer; -import org.apache.kafka.common.serialization.Deserializer; -import org.apache.kafka.common.serialization.IntegerDeserializer; -import org.apache.kafka.common.serialization.IntegerSerializer; -import org.apache.kafka.common.serialization.Serializer; -import org.apache.kafka.streaming.kstream.internals.FilteredIterator; -import org.apache.kafka.streaming.kstream.internals.WindowSupport; -import org.apache.kafka.streaming.processor.internals.ProcessorContextImpl; -import org.apache.kafka.streaming.processor.internals.RecordCollector; -import org.apache.kafka.streaming.processor.ProcessorContext; -import org.apache.kafka.streaming.processor.RestoreFunc; -import org.apache.kafka.streaming.processor.internals.Stamped; - -import java.util.HashMap; -import java.util.Iterator; -import java.util.LinkedList; -import java.util.Map; - -public class SlidingWindow extends WindowSupport implements Window { - - private final Object lock = new Object(); - private final Serializer keySerializer; - private final Serializer valueSerializer; - private final Deserializer keyDeserializer; - private final Deserializer valueDeserializer; - private ProcessorContext context; - private int slotNum; - private String name; - private final long duration; - private final int maxCount; - private LinkedList list = new LinkedList(); - private HashMap> map = new HashMap<>(); - - public SlidingWindow( - String name, - long duration, - int maxCount, - Serializer keySerializer, - Serializer valueSerializer, - Deserializer keyDeseriaizer, - Deserializer valueDeserializer) { - this.name = name; - this.duration = duration; - this.maxCount = maxCount; - this.keySerializer = keySerializer; - this.valueSerializer = valueSerializer; - this.keyDeserializer = keyDeseriaizer; - this.valueDeserializer = valueDeserializer; - } - - @Override - public void init(ProcessorContext context) { - this.context = context; - RestoreFuncImpl restoreFunc = new RestoreFuncImpl(); - context.register(this, restoreFunc); - - for (ValueList valueList : map.values()) { - valueList.clearDirtyValues(); - } - this.slotNum = restoreFunc.slotNum; - } - - @Override - public Iterator findAfter(K key, final long timestamp) { - return find(key, timestamp, timestamp + duration); - } - - @Override - public Iterator findBefore(K key, final long timestamp) { - return find(key, timestamp - duration, timestamp); - } - - @Override - public Iterator find(K key, final long timestamp) { - return find(key, timestamp - duration, timestamp + duration); - } - - /* - * finds items in the window between startTime and endTime (both inclusive) - */ - private Iterator find(K key, final long startTime, final long endTime) { - final ValueList values = map.get(key); - - if (values == null) { - return null; - } else { - return new FilteredIterator>(values.iterator()) { - @Override - protected V filter(Value item) { - if (startTime <= item.timestamp && item.timestamp <= endTime) - return item.value; - else - return null; - } - }; - } - } - - @Override - public void put(K key, V value, long timestamp) { - synchronized (lock) { - slotNum++; - - list.offerLast(key); - - ValueList values = map.get(key); - if (values == null) { - values = new ValueList<>(); - map.put(key, values); - } - - values.add(slotNum, value, timestamp); - } - evictExcess(); - evictExpired(timestamp - duration); - } - - private void evictExcess() { - while (list.size() > maxCount) { - K oldestKey = list.pollFirst(); - - ValueList values = map.get(oldestKey); - values.removeFirst(); - - if (values.isEmpty()) map.remove(oldestKey); - } - } - - private void evictExpired(long cutoffTime) { - while (true) { - K oldestKey = list.peekFirst(); - - ValueList values = map.get(oldestKey); - Stamped oldestValue = values.first(); - - if (oldestValue.timestamp < cutoffTime) { - list.pollFirst(); - values.removeFirst(); - - if (values.isEmpty()) map.remove(oldestKey); - } else { - break; - } - } - } - - @Override - public String name() { - return name; - } - - @Override - public void flush() { - IntegerSerializer intSerializer = new IntegerSerializer(); - ByteArraySerializer byteArraySerializer = new ByteArraySerializer(); - - RecordCollector collector = ((ProcessorContextImpl)context).recordCollector(); - - for (Map.Entry> entry : map.entrySet()) { - ValueList values = entry.getValue(); - if (values.hasDirtyValues()) { - K key = entry.getKey(); - - byte[] keyBytes = keySerializer.serialize(name, key); - - Iterator> iterator = values.dirtyValueIterator(); - while (iterator.hasNext()) { - Value dirtyValue = iterator.next(); - byte[] slot = intSerializer.serialize("", dirtyValue.slotNum); - byte[] valBytes = valueSerializer.serialize(name, dirtyValue.value); - - byte[] combined = new byte[8 + 4 + keyBytes.length + 4 + valBytes.length]; - - int offset = 0; - offset += putLong(combined, offset, dirtyValue.timestamp); - offset += puts(combined, offset, keyBytes); - offset += puts(combined, offset, valBytes); - - if (offset != combined.length) throw new IllegalStateException("serialized length does not match"); - - collector.send(new ProducerRecord<>(name, context.id(), slot, combined), byteArraySerializer, byteArraySerializer); - } - values.clearDirtyValues(); - } - } - } - - @Override - public void close() { - // TODO - } - - @Override - public boolean persistent() { - // TODO: should not be persistent, right? - return false; - } - - private class RestoreFuncImpl implements RestoreFunc { - - final IntegerDeserializer intDeserializer; - int slotNum = 0; - - RestoreFuncImpl() { - intDeserializer = new IntegerDeserializer(); - } - - @Override - public void apply(byte[] slot, byte[] bytes) { - - slotNum = intDeserializer.deserialize("", slot); - - int offset = 0; - // timestamp - long timestamp = getLong(bytes, offset); - offset += 8; - // key - int length = getInt(bytes, offset); - offset += 4; - K key = deserialize(bytes, offset, length, name, keyDeserializer); - offset += length; - // value - length = getInt(bytes, offset); - offset += 4; - V value = deserialize(bytes, offset, length, name, valueDeserializer); - - put(key, value, timestamp); - } - } - -} diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/SlidingWindowDef.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/SlidingWindowDef.java new file mode 100644 index 0000000000000..4fc8ca41411c8 --- /dev/null +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/SlidingWindowDef.java @@ -0,0 +1,264 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.kafka.streaming.kstream; + +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.serialization.ByteArraySerializer; +import org.apache.kafka.common.serialization.Deserializer; +import org.apache.kafka.common.serialization.IntegerDeserializer; +import org.apache.kafka.common.serialization.IntegerSerializer; +import org.apache.kafka.common.serialization.Serializer; +import org.apache.kafka.streaming.kstream.internals.FilteredIterator; +import org.apache.kafka.streaming.kstream.internals.WindowSupport; +import org.apache.kafka.streaming.processor.internals.ProcessorContextImpl; +import org.apache.kafka.streaming.processor.internals.RecordCollector; +import org.apache.kafka.streaming.processor.ProcessorContext; +import org.apache.kafka.streaming.processor.RestoreFunc; +import org.apache.kafka.streaming.processor.internals.Stamped; + +import java.util.HashMap; +import java.util.Iterator; +import java.util.LinkedList; +import java.util.Map; + +public class SlidingWindowDef implements WindowDef { + private String name; + private final long duration; + private final int maxCount; + private final Serializer keySerializer; + private final Serializer valueSerializer; + private final Deserializer keyDeserializer; + private final Deserializer valueDeserializer; + + public SlidingWindowDef( + String name, + long duration, + int maxCount, + Serializer keySerializer, + Serializer valueSerializer, + Deserializer keyDeseriaizer, + Deserializer valueDeserializer) { + this.name = name; + this.duration = duration; + this.maxCount = maxCount; + this.keySerializer = keySerializer; + this.valueSerializer = valueSerializer; + this.keyDeserializer = keyDeseriaizer; + this.valueDeserializer = valueDeserializer; + } + + @Override + public String name() { + return name; + } + + @Override + public Window build() { + return new SlidingWindow(); + } + + public class SlidingWindow extends WindowSupport implements Window { + private final Object lock = new Object(); + private ProcessorContext context; + private int slotNum; // used as a key for Kafka log compaction + private LinkedList list = new LinkedList(); + private HashMap> map = new HashMap<>(); + + @Override + public void init(ProcessorContext context) { + this.context = context; + RestoreFuncImpl restoreFunc = new RestoreFuncImpl(); + context.register(this, restoreFunc); + + for (ValueList valueList : map.values()) { + valueList.clearDirtyValues(); + } + this.slotNum = restoreFunc.slotNum; + } + + @Override + public Iterator findAfter(K key, final long timestamp) { + return find(key, timestamp, timestamp + duration); + } + + @Override + public Iterator findBefore(K key, final long timestamp) { + return find(key, timestamp - duration, timestamp); + } + + @Override + public Iterator find(K key, final long timestamp) { + return find(key, timestamp - duration, timestamp + duration); + } + + /* + * finds items in the window between startTime and endTime (both inclusive) + */ + private Iterator find(K key, final long startTime, final long endTime) { + final ValueList values = map.get(key); + + if (values == null) { + return null; + } else { + return new FilteredIterator>(values.iterator()) { + @Override + protected V filter(Value item) { + if (startTime <= item.timestamp && item.timestamp <= endTime) + return item.value; + else + return null; + } + }; + } + } + + @Override + public void put(K key, V value, long timestamp) { + synchronized (lock) { + slotNum++; + + list.offerLast(key); + + ValueList values = map.get(key); + if (values == null) { + values = new ValueList<>(); + map.put(key, values); + } + + values.add(slotNum, value, timestamp); + } + evictExcess(); + evictExpired(timestamp - duration); + } + + private void evictExcess() { + while (list.size() > maxCount) { + K oldestKey = list.pollFirst(); + + ValueList values = map.get(oldestKey); + values.removeFirst(); + + if (values.isEmpty()) map.remove(oldestKey); + } + } + + private void evictExpired(long cutoffTime) { + while (true) { + K oldestKey = list.peekFirst(); + + ValueList values = map.get(oldestKey); + Stamped oldestValue = values.first(); + + if (oldestValue.timestamp < cutoffTime) { + list.pollFirst(); + values.removeFirst(); + + if (values.isEmpty()) map.remove(oldestKey); + } else { + break; + } + } + } + + @Override + public String name() { + return name; + } + + @Override + public void flush() { + IntegerSerializer intSerializer = new IntegerSerializer(); + ByteArraySerializer byteArraySerializer = new ByteArraySerializer(); + + RecordCollector collector = ((ProcessorContextImpl) context).recordCollector(); + + for (Map.Entry> entry : map.entrySet()) { + ValueList values = entry.getValue(); + if (values.hasDirtyValues()) { + K key = entry.getKey(); + + byte[] keyBytes = keySerializer.serialize(name, key); + + Iterator> iterator = values.dirtyValueIterator(); + while (iterator.hasNext()) { + Value dirtyValue = iterator.next(); + byte[] slot = intSerializer.serialize("", dirtyValue.slotNum); + byte[] valBytes = valueSerializer.serialize(name, dirtyValue.value); + + byte[] combined = new byte[8 + 4 + keyBytes.length + 4 + valBytes.length]; + + int offset = 0; + offset += putLong(combined, offset, dirtyValue.timestamp); + offset += puts(combined, offset, keyBytes); + offset += puts(combined, offset, valBytes); + + if (offset != combined.length) + throw new IllegalStateException("serialized length does not match"); + + collector.send(new ProducerRecord<>(name, context.id(), slot, combined), byteArraySerializer, byteArraySerializer); + } + values.clearDirtyValues(); + } + } + } + + @Override + public void close() { + // TODO + } + + @Override + public boolean persistent() { + // TODO: should not be persistent, right? + return false; + } + + private class RestoreFuncImpl implements RestoreFunc { + + final IntegerDeserializer intDeserializer; + int slotNum = 0; + + RestoreFuncImpl() { + intDeserializer = new IntegerDeserializer(); + } + + @Override + public void apply(byte[] slot, byte[] bytes) { + + slotNum = intDeserializer.deserialize("", slot); + + int offset = 0; + // timestamp + long timestamp = getLong(bytes, offset); + offset += 8; + // key + int length = getInt(bytes, offset); + offset += 4; + K key = deserialize(bytes, offset, length, name, keyDeserializer); + offset += length; + // value + length = getInt(bytes, offset); + offset += 4; + V value = deserialize(bytes, offset, length, name, valueDeserializer); + + put(key, value, timestamp); + } + } + } + +} diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/Window.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/Window.java index d805466144d4c..fc9582d21fc20 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/Window.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/Window.java @@ -33,5 +33,4 @@ public interface Window extends StateStore { Iterator findBefore(K key, long timestamp); void put(K key, V value, long timestamp); - } diff --git a/stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorMetadata.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/WindowDef.java similarity index 72% rename from stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorMetadata.java rename to stream/src/main/java/org/apache/kafka/streaming/kstream/WindowDef.java index ba773b0f279ca..26971937eccfe 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorMetadata.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/WindowDef.java @@ -15,19 +15,11 @@ * limitations under the License. */ -package org.apache.kafka.streaming.processor; +package org.apache.kafka.streaming.kstream; -public class ProcessorMetadata { +public interface WindowDef { - private String name; - private Object value; + String name(); - public ProcessorMetadata(String name, Object value) { - this.name = name; - this.value = value; - } - - public Object value() { - return value; - } + Window build(); } 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 25189f677c994..e9d6030401c58 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 @@ -17,36 +17,34 @@ package org.apache.kafka.streaming.kstream.internals; -import org.apache.kafka.common.KafkaException; import org.apache.kafka.streaming.processor.Processor; -import org.apache.kafka.streaming.processor.ProcessorMetadata; +import org.apache.kafka.streaming.processor.ProcessorFactory; import org.apache.kafka.streaming.kstream.Predicate; -class KStreamBranch extends Processor { +class KStreamBranch implements ProcessorFactory { private final Predicate[] predicates; @SuppressWarnings("unchecked") - public KStreamBranch(String name, ProcessorMetadata config) { - super(name, config); - - if (this.metadata() == null) - throw new IllegalStateException("ProcessorMetadata should be specified."); - - this.predicates = (Predicate[]) config.value(); + public KStreamBranch(Predicate... predicates) { + this.predicates = predicates; } @Override - public void process(K key, V value) { - if (this.children().size() != this.predicates.length) - throw new KafkaException("Number of branched streams does not match the length of predicates: this should not happen."); + public Processor build() { + return new KStreamBranchProcessor(); + } - for (int i = 0; i < predicates.length; i++) { - if (predicates[i].apply(key, value)) { - // do not use forward here bu directly call process() and then break the loop - // so that no record is going to be piped to multiple streams - this.children().get(i).process(key, value); - break; + private class KStreamBranchProcessor extends KStreamProcessor { + @Override + public void process(K key, V value) { + for (int i = 0; i < predicates.length; i++) { + if (predicates[i].apply(key, value)) { + // use forward with childIndex here and then break the loop + // so that no record is going to be piped to multiple streams + context.forward(key, value, i); + break; + } } } } 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 b467a9d3d0fc3..a85e0181a7c76 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,42 +19,29 @@ import org.apache.kafka.streaming.processor.Processor; import org.apache.kafka.streaming.kstream.Predicate; -import org.apache.kafka.streaming.processor.ProcessorMetadata; +import org.apache.kafka.streaming.processor.ProcessorFactory; -class KStreamFilter extends Processor { +class KStreamFilter implements ProcessorFactory { - private final PredicateOut predicateOut; + private final Predicate predicate; + private final boolean filterOut; - public static final class PredicateOut { - - public final Predicate predicate; - public final boolean filterOut; - - public PredicateOut(Predicate predicate) { - this(predicate, false); - } - - public PredicateOut(Predicate predicate, boolean filterOut) { - this.predicate = predicate; - this.filterOut = filterOut; - } + public KStreamFilter(Predicate predicate, boolean filterOut) { + this.predicate = predicate; + this.filterOut = filterOut; } - @SuppressWarnings("unchecked") - public KStreamFilter(ProcessorMetadata metadata) { - super(metadata); - - if (this.metadata() == null) - throw new IllegalStateException("ProcessorMetadata should be specified."); - - this.predicateOut = (PredicateOut) metadata.value(); + @Override + public Processor build() { + return new KStreamFilterProcessor(); } - @Override - public void process(K key, V value) { - if ((!predicateOut.filterOut && predicateOut.predicate.apply(key, value)) - || (predicateOut.filterOut && !predicateOut.predicate.apply(key, value))) { - context.forward(key, value); + private class KStreamFilterProcessor extends KStreamProcessor { + @Override + public void process(K key, V value) { + if (filterOut ^ predicate.apply(key, value)) { + context.forward(key, value); + } } } } 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 89543b704e986..bc834a8f66ae8 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 @@ -17,30 +17,30 @@ package org.apache.kafka.streaming.kstream.internals; -import org.apache.kafka.streaming.processor.Processor; -import org.apache.kafka.streaming.processor.ProcessorMetadata; 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; -class KStreamFlatMap extends Processor { +class KStreamFlatMap implements ProcessorFactory { private final KeyValueFlatMap mapper; - @SuppressWarnings("unchecked") - KStreamFlatMap(ProcessorMetadata metadata) { - super(metadata); - - if (this.metadata() == null) - throw new IllegalStateException("ProcessorMetadata should be specified."); - - this.mapper = (KeyValueFlatMap) metadata.value(); + KStreamFlatMap(KeyValueFlatMap mapper) { + this.mapper = mapper; } @Override - public void process(K1 key, V1 value) { - Iterable> pairs = mapper.apply(key, value); - for (KeyValue pair : pairs) { - context.forward(pair.key, pair.value); + public Processor build() { + return new KStreamFlatMapProcessor(); + } + + private class KStreamFlatMapProcessor extends KStreamProcessor { + @Override + public void process(K1 key, V1 value) { + for (KeyValue newPair : mapper.apply(key, value)) { + context.forward(newPair.key, newPair.value); + } } } } 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 0df01e62358ed..5e18f188ffa08 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,27 +19,29 @@ import org.apache.kafka.streaming.processor.Processor; import org.apache.kafka.streaming.kstream.ValueMapper; -import org.apache.kafka.streaming.processor.ProcessorMetadata; +import org.apache.kafka.streaming.processor.ProcessorFactory; -class KStreamFlatMapValues extends Processor { +class KStreamFlatMapValues implements ProcessorFactory { private final ValueMapper> mapper; @SuppressWarnings("unchecked") - KStreamFlatMapValues(String name, ProcessorMetadata config) { - super(name, config); - - if (this.metadata() == null) - throw new IllegalStateException("ProcessorMetadata should be specified."); - - this.mapper = (ValueMapper>) config.value(); + KStreamFlatMapValues(ValueMapper> mapper) { + this.mapper = mapper; } @Override - public void process(K1 key, V1 value) { - Iterable newValues = mapper.apply(value); - for (V2 v : newValues) { - forward(key, v); + public Processor build() { + return new KStreamFlatMapValuesProcessor(); + } + + private class KStreamFlatMapValuesProcessor extends KStreamProcessor { + @Override + public void process(K1 key, V1 value) { + Iterable newValues = mapper.apply(value); + for (V2 v : newValues) { + context.forward(key, v); + } } } } 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 14b09c52bceb3..e317338b134e6 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,17 +20,16 @@ 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.ProcessorMetadata; +import org.apache.kafka.streaming.processor.ProcessorFactory; import org.apache.kafka.streaming.processor.TopologyBuilder; import org.apache.kafka.streaming.kstream.KStreamWindowed; import org.apache.kafka.streaming.kstream.KeyValueMapper; import org.apache.kafka.streaming.kstream.Predicate; import org.apache.kafka.streaming.kstream.KStream; import org.apache.kafka.streaming.kstream.ValueMapper; -import org.apache.kafka.streaming.kstream.Window; +import org.apache.kafka.streaming.kstream.WindowDef; import java.util.ArrayList; -import java.util.Arrays; import java.util.List; import java.util.concurrent.atomic.AtomicInteger; @@ -50,10 +49,20 @@ public class KStreamImpl implements KStream { private static final String BRANCH_NAME = "KAFKA-BRANCH-"; + private static final String BRANCHCHILD_NAME = "KAFKA-BRANCHCHILD-"; + + private static final String WINDOWED_NAME = "KAFKA-WINDOWED-"; + + public static final String JOIN_NAME = "KAFKA-JOIN-"; + + public static final String JOINOTHER_NAME = "KAFKA-JOINOTHER-"; + public static final String SOURCE_NAME = "KAFKA-SOURCE-"; public static final String SEND_NAME = "KAFKA-SEND-"; + public static final String WINDOW_NAME = "KAFKA-WINDOW-"; + public static final AtomicInteger INDEX = new AtomicInteger(1); protected TopologyBuilder topology; @@ -68,7 +77,7 @@ public KStreamImpl(TopologyBuilder topology, String name) { public KStream filter(Predicate predicate) { String name = FILTER_NAME + INDEX.getAndIncrement(); - topology.addProcessor(name, KStreamFilter.class, new ProcessorMetadata("Predicate", new KStreamFilter.PredicateOut<>(predicate)), this.name); + topology.addProcessor(name, new KStreamFilter(predicate, false), this.name); return new KStreamImpl<>(topology, name); } @@ -77,7 +86,7 @@ public KStream filter(Predicate predicate) { public KStream filterOut(final Predicate predicate) { String name = FILTER_NAME + INDEX.getAndIncrement(); - topology.addProcessor(name, KStreamFilter.class, new ProcessorMetadata("Predicate", new KStreamFilter.PredicateOut<>(predicate, true)), this.name); + topology.addProcessor(name, new KStreamFilter(predicate, true), this.name); return new KStreamImpl<>(topology, name); } @@ -86,7 +95,7 @@ public KStream filterOut(final Predicate predicate) { public KStream map(KeyValueMapper mapper) { String name = MAP_NAME + INDEX.getAndIncrement(); - topology.addProcessor(name, KStreamMap.class, new ProcessorMetadata("Mapper", mapper), this.name); + topology.addProcessor(name, new KStreamMap(mapper), this.name); return new KStreamImpl<>(topology, name); } @@ -95,7 +104,7 @@ public KStream map(KeyValueMapper mapper) { public KStream mapValues(ValueMapper mapper) { String name = MAPVALUES_NAME + INDEX.getAndIncrement(); - topology.addProcessor(name, KStreamMapValues.class, new ProcessorMetadata("ValueMapper", mapper), this.name); + topology.addProcessor(name, new KStreamMapValues(mapper), this.name); return new KStreamImpl<>(topology, name); } @@ -105,7 +114,7 @@ public KStream mapValues(ValueMapper mapper) { public KStream flatMap(KeyValueFlatMap mapper) { String name = FLATMAP_NAME + INDEX.getAndIncrement(); - topology.addProcessor(name, KStreamFlatMap.class, new ProcessorMetadata("FlatMapper", mapper), this.name); + topology.addProcessor(name, new KStreamFlatMap(mapper), this.name); return new KStreamImpl<>(topology, name); } @@ -115,34 +124,37 @@ public KStream flatMap(KeyValueFlatMap mapper) { public KStream flatMapValues(ValueMapper> mapper) { String name = FLATMAPVALUES_NAME + INDEX.getAndIncrement(); - topology.addProcessor(name, KStreamFlatMapValues.class, new ProcessorMetadata("ValueMapper", mapper), this.name); + topology.addProcessor(name, new KStreamFlatMapValues(mapper), this.name); return new KStreamImpl<>(topology, name); } @Override - public KStreamWindowed with(Window window) { - KStreamWindow windowed = new KStreamWindow<>(window); + public KStreamWindowed with(WindowDef windowDef) { + String name = WINDOWED_NAME + INDEX.getAndIncrement(); - topology.addProcessor(windowed, processor); + topology.addProcessor(name, new KStreamWindow<>(windowDef), this.name); - return new KStreamWindow.KStreamWindowedImpl<>(topology, windowed); + return new KStreamWindowedImpl<>(topology, name, windowDef); } @Override @SuppressWarnings("unchecked") public KStream[] branch(Predicate... predicates) { - String name = BRANCH_NAME + INDEX.getAndIncrement(); + String branchName = BRANCH_NAME + INDEX.getAndIncrement(); - topology.addProcessor(name, KStreamBranch.class, new ProcessorMetadata("Predicates", Arrays.copyOf(predicates, predicates.length)), this.name); + topology.addProcessor(branchName, new KStreamBranch(predicates.clone()), this.name); - KStreamImpl branch = new KStreamImpl<>(topology, name); - List> avatars = new ArrayList<>(); + List> branchChildren = new ArrayList<>(); for (int i = 0; i < predicates.length; i++) { - avatars.add(branch); + String childName = BRANCHCHILD_NAME + INDEX.getAndIncrement(); + + topology.addProcessor(childName, new KStreamPassThrough(), branchName); + + branchChildren.add(new KStreamImpl(topology, childName)); } - return (KStream[]) avatars.toArray(); + return (KStream[]) branchChildren.toArray(); } @SuppressWarnings("unchecked") @@ -154,7 +166,7 @@ public KStream through(String topic, Deserializer valDeserializer) { String sendName = SEND_NAME + INDEX.getAndIncrement(); - process(new KStreamSend(sendName, new ProcessorMetadata("Topic-Ser", new KStreamSend.TopicSer(topic, (Serializer) keySerializer, (Serializer) valSerializer)))); + topology.addSink(sendName, topic, keySerializer, valSerializer, this.name); String sourceName = SOURCE_NAME + INDEX.getAndIncrement(); @@ -168,15 +180,15 @@ public KStream through(String topic, public void sendTo(String topic, Serializer keySerializer, Serializer valSerializer) { String name = SEND_NAME + INDEX.getAndIncrement(); - process(new KStreamSend(name, new ProcessorMetadata("Topic-Ser", new KStreamSend.TopicSer(topic, (Serializer) keySerializer, (Serializer) valSerializer)))); + topology.addSink(name, topic, keySerializer, valSerializer, this.name); } @SuppressWarnings("unchecked") @Override - public KStream process(final Processor processor) { + public KStream process(final ProcessorFactory processorFactory) { String name = PROCESSOR_NAME + INDEX.getAndIncrement(); - topology.addProcessor(name, KStreamProcessor.class, new ProcessorMetadata("Processor", processor), this.name); + topology.addProcessor(name, processorFactory, 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 7751a7f3603f3..a8452d5ff0b64 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,109 +21,111 @@ 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 java.util.Iterator; -class KStreamJoin extends Processor { - - private static final String JOIN_NAME = "KAFKA-JOIN"; - private static final String JOIN_OTHER_NAME = "KAFKA-JOIN-OTHER"; +class KStreamJoin implements ProcessorFactory { private static abstract class Finder { abstract Iterator find(K key, long timestamp); } - private final KStreamWindow stream1; - private final KStreamWindow stream2; - private final Finder finder1; - private final Finder finder2; + private final String windowName1; + private final String windowName2; private final ValueJoiner joiner; - final Processor processorForOtherStream; - - private ProcessorContext context; - - KStreamJoin(KStreamWindow stream1, KStreamWindow stream2, boolean prior, ValueJoiner joiner) { - super(JOIN_NAME); + private final boolean prior; - this.stream1 = stream1; - this.stream2 = stream2; - final Window window1 = stream1.window(); - final Window window2 = stream2.window(); - - if (prior) { - this.finder1 = new Finder() { - Iterator find(K key, long timestamp) { - return window1.findAfter(key, timestamp); - } - }; - this.finder2 = new Finder() { - Iterator find(K key, long timestamp) { - return window2.findBefore(key, timestamp); - } - }; - } else { - this.finder1 = new Finder() { - Iterator find(K key, long timestamp) { - return window1.find(key, timestamp); - } - }; - this.finder2 = new Finder() { - Iterator find(K key, long timestamp) { - return window2.find(key, timestamp); - } - }; + private Processor processorForOtherStream = null; + public final ProcessorFactory processorFactoryForOtherStream = new ProcessorFactory() { + @Override + public Processor build() { + return processorForOtherStream; } + }; + KStreamJoin(String windowName1, String windowName2, boolean prior, ValueJoiner joiner) { + this.windowName1 = windowName1; + this.windowName2 = windowName2; this.joiner = joiner; - - this.processorForOtherStream = processorForOther(); + this.prior = prior; } @Override - public void init(ProcessorContext context) { - this.context = context; - - // check if these two streams are joinable - if (!stream1.context().joinable(stream2.context())) - throw new IllegalStateException("Stream " + stream1.name() + " and stream " + - stream2.name() + " are not joinable."); + public Processor build() { + return new KStreamJoinProcessor(); } - @Override - public void process(K key, V1 value) { - long timestamp = context.timestamp(); - Iterator iter = finder2.find(key, timestamp); - if (iter != null) { - while (iter.hasNext()) { - doJoin(key, value, iter.next()); + private class KStreamJoinProcessor extends KStreamProcessor { + + private Finder finder1; + private Finder finder2; + + @SuppressWarnings("unchecked") + @Override + public void init(ProcessorContext context) { + super.init(context); + + // check if these two streams are joinable + if (!context.joinable()) + throw new IllegalStateException("Streams are not joinable."); + + final Window window1 = (Window) context.getStateStore(windowName1); + final Window window2 = (Window) context.getStateStore(windowName2); + + if (prior) { + this.finder1 = new Finder() { + Iterator find(K key, long timestamp) { + return window1.findAfter(key, timestamp); + } + }; + this.finder2 = new Finder() { + Iterator find(K key, long timestamp) { + return window2.findBefore(key, timestamp); + } + }; + } else { + this.finder1 = new Finder() { + Iterator find(K key, long timestamp) { + return window1.find(key, timestamp); + } + }; + this.finder2 = new Finder() { + Iterator find(K key, long timestamp) { + return window2.find(key, timestamp); + } + }; } - } - } - private Processor processorForOther() { - return new Processor(JOIN_OTHER_NAME) { - - @SuppressWarnings("unchecked") - @Override - public void process(K key, V2 value) { - long timestamp = context.timestamp(); - Iterator iter = finder1.find(key, timestamp); - if (iter != null) { - while (iter.hasNext()) { - doJoin(key, iter.next(), value); + processorForOtherStream = new KStreamProcessor() { + @Override + public void process(K key, V2 value) { + long timestamp = context.timestamp(); + Iterator iter = finder1.find(key, timestamp); + if (iter != null) { + while (iter.hasNext()) { + doJoin(key, iter.next(), value); + } } } - } + }; + } - @Override - public void close() { - // down stream instances are close when the primary stream is closed + @Override + public void process(K key, V1 value) { + long timestamp = context.timestamp(); + Iterator iter = finder2.find(key, timestamp); + if (iter != null) { + while (iter.hasNext()) { + doJoin(key, value, iter.next()); + } } - }; - } + } - // TODO: use the "outer-stream" topic as the resulted join stream topic - private void doJoin(K key, V1 value1, V2 value2) { - forward(key, joiner.apply(value1, value2)); + // TODO: use the "outer-stream" topic as the resulted join stream topic + private void doJoin(K key, V1 value1, V2 value2) { + context.forward(key, joiner.apply(value1, value2)); + } } + } 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 903fb785a7561..83fb4c11c99f6 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,25 +20,26 @@ 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.ProcessorMetadata; +import org.apache.kafka.streaming.processor.ProcessorFactory; -class KStreamMap extends Processor { +class KStreamMap implements ProcessorFactory { private final KeyValueMapper mapper; - @SuppressWarnings("unchecked") - public KStreamMap(ProcessorMetadata metadata) { - super(metadata); - - if (this.metadata() == null) - throw new IllegalStateException("ProcessorMetadata should be specified."); - - this.mapper = (KeyValueMapper) metadata.value(); + public KStreamMap(KeyValueMapper mapper) { + this.mapper = mapper; } @Override - public void process(K1 key, V1 value) { - KeyValue newPair = mapper.apply(key, value); - context.forward(newPair.key, newPair.value); + public Processor build() { + return new KStreamMapProcessor(); + } + + private class KStreamMapProcessor extends KStreamProcessor { + @Override + public void process(K1 key, V1 value) { + KeyValue newPair = mapper.apply(key, value); + context.forward(newPair.key, newPair.value); + } } } 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 a8aac1285039e..667c04386c381 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,25 +19,26 @@ import org.apache.kafka.streaming.processor.Processor; import org.apache.kafka.streaming.kstream.ValueMapper; -import org.apache.kafka.streaming.processor.ProcessorMetadata; +import org.apache.kafka.streaming.processor.ProcessorFactory; -class KStreamMapValues extends Processor { +class KStreamMapValues implements ProcessorFactory { private final ValueMapper mapper; - @SuppressWarnings("unchecked") - public KStreamMapValues(String name, ProcessorMetadata config) { - super(name, config); - - if (this.metadata() == null) - throw new IllegalStateException("ProcessorMetadata should be specified."); - - this.mapper = (ValueMapper) config.value(); + public KStreamMapValues(ValueMapper mapper) { + this.mapper = mapper; } @Override - public void process(K1 key, V1 value) { - V2 newValue = mapper.apply(value); - forward(key, newValue); + public Processor build() { + return new KStreamMapProcessor(); + } + + private class KStreamMapProcessor extends KStreamProcessor { + @Override + public void process(K1 key, V1 value) { + V2 newValue = mapper.apply(value); + context.forward(key, newValue); + } } } 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 new file mode 100644 index 0000000000000..98528f72c7c1b --- /dev/null +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamPassThrough.java @@ -0,0 +1,36 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.kafka.streaming.kstream.internals; + +import org.apache.kafka.streaming.processor.Processor; +import org.apache.kafka.streaming.processor.ProcessorFactory; + +class KStreamPassThrough implements ProcessorFactory { + + @Override + public Processor build() { + return new KStreamPassThroughProcessor(); + } + + public class KStreamPassThroughProcessor extends KStreamProcessor { + @Override + public void process(K key, V value) { + context.forward(key, value); + } + } +} diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamProcessor.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamProcessor.java index 486b67e3c7625..a289893ae1fd0 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamProcessor.java +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamProcessor.java @@ -18,24 +18,27 @@ package org.apache.kafka.streaming.kstream.internals; import org.apache.kafka.streaming.processor.Processor; -import org.apache.kafka.streaming.processor.ProcessorMetadata; +import org.apache.kafka.streaming.processor.ProcessorContext; -public class KStreamProcessor extends Processor { +abstract class KStreamProcessor implements Processor { - private final Processor processor; + protected ProcessorContext context; - @SuppressWarnings("unchecked") - public KStreamProcessor(ProcessorMetadata metadata) { - super(metadata); + @Override + abstract public void process(K key, V value); - if (this.metadata() != null) - throw new IllegalStateException("ProcessorMetadata should be null."); + @Override + public void init(ProcessorContext context) { + this.context = context; + } - this.processor = (Processor) metadata.value(); + @Override + public void punctuate(long streamTime) { + // do nothing } @Override - public void process(K key, V value) { - processor.process(key, value); + public void close() { + // do nothing } } diff --git a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamSend.java b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamSend.java deleted file mode 100644 index 12b5a8ffa2837..0000000000000 --- a/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamSend.java +++ /dev/null @@ -1,62 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.kafka.streaming.kstream.internals; - -import org.apache.kafka.common.serialization.Serializer; -import org.apache.kafka.streaming.processor.Processor; -import org.apache.kafka.streaming.processor.ProcessorMetadata; -import org.apache.kafka.streaming.processor.ProcessorContext; - -class KStreamSend extends Processor { - - private ProcessorContext context; - - private TopicSer topicSerializers; - - public static final class TopicSer { - public String topic; - public Serializer keySerializer; - public Serializer valSerializer; - - public TopicSer(String topic, Serializer keySerializer, Serializer valSerializer) { - this.topic = topic; - this.keySerializer = keySerializer; - this.valSerializer = valSerializer; - } - } - - @SuppressWarnings("unchecked") - public KStreamSend(String name, ProcessorMetadata config) { - super(name, config); - - if (this.metadata() == null) - throw new IllegalStateException("ProcessorMetadata should be specified."); - - this.topicSerializers = (TopicSer) config.value(); - } - - @Override - public void init(ProcessorContext context) { - this.context = context; - } - - @Override - public void process(K key, V value) { - this.context.send(topicSerializers.topic, key, value, topicSerializers.keySerializer, topicSerializers.valSerializer); - } -} \ No newline at end of file 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 2a6628ab5859a..f357613d94ce4 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 @@ -17,82 +17,52 @@ package org.apache.kafka.streaming.kstream.internals; +import org.apache.kafka.streaming.kstream.Window; import org.apache.kafka.streaming.processor.Processor; -import org.apache.kafka.streaming.processor.TopologyBuilder; +import org.apache.kafka.streaming.processor.ProcessorFactory; import org.apache.kafka.streaming.processor.ProcessorContext; -import org.apache.kafka.streaming.kstream.KStream; -import org.apache.kafka.streaming.kstream.KStreamWindowed; -import org.apache.kafka.streaming.kstream.ValueJoiner; -import org.apache.kafka.streaming.kstream.Window; - -public class KStreamWindow extends Processor { - - public static final class KStreamWindowedImpl extends KStreamImpl implements KStreamWindowed { +import org.apache.kafka.streaming.kstream.WindowDef; - public KStreamWindow windowed; +public class KStreamWindow implements ProcessorFactory { - public KStreamWindowedImpl(TopologyBuilder topology, KStreamWindow windowed) { - super(topology, windowed); - this.windowed = windowed; - } - - @Override - public KStream join(KStreamWindowed other, ValueJoiner processor) { - return join(other, false, processor); - } + private final WindowDef windowDef; - @Override - public KStream joinPrior(KStreamWindowed other, ValueJoiner processor) { - return join(other, true, processor); - } - - private KStream join(KStreamWindowed other, boolean prior, ValueJoiner processor) { - KStreamWindow thisWindow = this.windowed; - KStreamWindow otherWindow = ((KStreamWindowedImpl) other).windowed; - - KStreamJoin join = new KStreamJoin<>(thisWindow, otherWindow, prior, processor); - - topology.addProcessor(join, thisWindow); - topology.addProcessor(join.processorForOtherStream, otherWindow); - - return new KStreamImpl<>(topology, join); - } + KStreamWindow(WindowDef windowDef) { + this.windowDef = windowDef; } - private static final String WINDOW_NAME = "KAFKA-WINDOW"; - - private final Window window; - private ProcessorContext context; - - KStreamWindow(Window window) { - super(WINDOW_NAME); - this.window = window; + public WindowDef window() { + return windowDef; } - public Window window() { - return window; + @Override + public Processor build() { + return new KStreamWindowProcessor(); } - public ProcessorContext context() { - return context; - } + private class KStreamWindowProcessor extends KStreamProcessor { - @Override - public void init(ProcessorContext context) { - this.context = context; - } + private Window window; - @SuppressWarnings("unchecked") - @Override - public void process(K key, V value) { - synchronized (this) { - window.put(key, value, context.timestamp()); - forward(key, value); + @Override + public void init(ProcessorContext context) { + super.init(context); + this.window = windowDef.build(); + this.window.init(context); } - } - @Override - public void close() { - window.close(); + @SuppressWarnings("unchecked") + @Override + public void process(K key, V value) { + synchronized (this) { + window.put(key, value, context.timestamp()); + context.forward(key, value); + } + } + + @Override + public void close() { + window.close(); + } } } 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 new file mode 100644 index 0000000000000..e269c3ad40bcc --- /dev/null +++ b/stream/src/main/java/org/apache/kafka/streaming/kstream/internals/KStreamWindowedImpl.java @@ -0,0 +1,44 @@ +package org.apache.kafka.streaming.kstream.internals; + +import org.apache.kafka.streaming.kstream.KStream; +import org.apache.kafka.streaming.kstream.KStreamWindowed; +import org.apache.kafka.streaming.kstream.ValueJoiner; +import org.apache.kafka.streaming.kstream.WindowDef; +import org.apache.kafka.streaming.processor.TopologyBuilder; + +/** + * Created by yasuhiro on 8/27/15. + */ +public final class KStreamWindowedImpl extends KStreamImpl implements KStreamWindowed { + + private final WindowDef windowDef; + + public KStreamWindowedImpl(TopologyBuilder topology, String name, WindowDef windowDef) { + super(topology, name); + this.windowDef = windowDef; + } + + @Override + public KStream join(KStreamWindowed other, ValueJoiner valueJoiner) { + return join(other, false, valueJoiner); + } + + @Override + public KStream joinPrior(KStreamWindowed other, ValueJoiner valueJoiner) { + return join(other, true, valueJoiner); + } + + private KStream join(KStreamWindowed other, boolean prior, ValueJoiner valueJoiner) { + String thisWindowName = this.windowDef.name(); + String otherWindowName = ((KStreamWindowedImpl) other).windowDef.name(); + + KStreamJoin join = new KStreamJoin<>(thisWindowName, otherWindowName, prior, valueJoiner); + + 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); + + return new KStreamImpl<>(topology, joinName); + } +} diff --git a/stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorContext.java b/stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorContext.java index 15d2e2981eee7..db78fe52e365c 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorContext.java +++ b/stream/src/main/java/org/apache/kafka/streaming/processor/ProcessorContext.java @@ -26,7 +26,7 @@ public interface ProcessorContext { // TODO: this is better moved to a KStreamContext - boolean joinable(ProcessorContext other); + boolean joinable(); /** * Returns the partition group id @@ -84,10 +84,14 @@ public interface ProcessorContext { */ void register(StateStore store, RestoreFunc restoreFunc); + StateStore getStateStore(String name); + void schedule(Processor processor, long interval); void forward(K key, V value); + void forward(K key, V value, int childIndex); + void commit(); String topic(); 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 6ee522a0dc0f3..894e0e37f0244 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 @@ -172,7 +172,7 @@ public ProcessorTopology build() { if (factory instanceof ProcessorNodeFactory) { for (String parent : ((ProcessorNodeFactory) factory).parents) { - processorMap.get(parent).chain(node); + processorMap.get(parent).addChild(node); } } else if (factory instanceof SourceNodeFactory) { for (String topic : ((SourceNodeFactory)factory).topics) { @@ -183,7 +183,7 @@ public ProcessorTopology build() { topicSinkMap.put(topic, (SinkNode) node); for (String parent : ((SinkNodeFactory) factory).parents) { - processorMap.get(parent).chain(node); + processorMap.get(parent).addChild(node); } } else { throw new IllegalStateException("unknown factory class: " + factory.getClass().getName()); diff --git a/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorContextImpl.java b/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorContextImpl.java index 1691632d879a0..fcff4b0c1b9b8 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorContextImpl.java +++ b/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorContextImpl.java @@ -94,13 +94,7 @@ public RecordCollector recordCollector() { } @Override - public boolean joinable(ProcessorContext o) { - - ProcessorContextImpl other = (ProcessorContextImpl) o; - - if (this.task != other.task) - return false; - + public boolean joinable() { Set partitions = this.task.partitions(); Map> partitionsById = new HashMap<>(); int firstId = -1; @@ -171,6 +165,10 @@ public void register(StateStore store, RestoreFunc restoreFunc) { stateMgr.register(store, restoreFunc); } + public StateStore getStateStore(String name) { + return stateMgr.getStore(name); + } + @Override public String topic() { if (task.record() == null) @@ -206,12 +204,20 @@ public long timestamp() { @Override @SuppressWarnings("unchecked") public void forward(K key, V value) { - for (ProcessorNode childNode : (List>) task.node().children()) { + for (ProcessorNode childNode : (List>) task.node().children()) { task.node(childNode); childNode.process(key, value); } } + @Override + @SuppressWarnings("unchecked") + public void forward(K key, V value, int childIndex) { + ProcessorNode childNode = (ProcessorNode) task.node().children().get(childIndex); + task.node(childNode); + childNode.process(key, value); + } + @Override public void commit() { task.commitOffset(); 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 480c577754fdb..fa1dbb1e734bc 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 @@ -23,22 +23,20 @@ import java.util.ArrayList; import java.util.List; -public class ProcessorNode { +public class ProcessorNode { - private final List> children; - private final List> parents; + private final List> children; private final String name; - private final Processor processor; + private final Processor processor; public ProcessorNode(String name) { this(name, null); } - public ProcessorNode(String name, Processor processor) { + public ProcessorNode(String name, Processor processor) { this.name = name; this.processor = processor; - this.parents = new ArrayList<>(); this.children = new ArrayList<>(); } @@ -46,20 +44,11 @@ public String name() { return name; } - public Processor processor() { - return processor; - } - - public List> parents() { - return parents; - } - - public List> children() { + public List> children() { return children; } - public final void chain(ProcessorNode child) { - child.parents.add(this); + public final void addChild(ProcessorNode child) { children.add(child); } @@ -67,7 +56,7 @@ public void init(ProcessorContext context) { processor.init(context); } - public void process(K1 key, V1 value) { + public void process(K key, V value) { processor.process(key, value); } diff --git a/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorStateManager.java b/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorStateManager.java index 5052421e0fe9f..c659e0459d2e2 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorStateManager.java +++ b/stream/src/main/java/org/apache/kafka/streaming/processor/internals/ProcessorStateManager.java @@ -77,7 +77,7 @@ public void register(StateStore store, RestoreFunc restoreFunc) { // ---- register the store ---- // // check that the underlying change log topic exist or not - if (restoreConsumer.listTopics().keySet().contains(store.name())) { + if (restoreConsumer.listTopics().containsKey(store.name())) { boolean partitionNotFound = true; for (PartitionInfo partitionInfo : restoreConsumer.partitionsFor(store.name())) { if (partitionInfo.partition() == id) { @@ -139,6 +139,10 @@ public void register(StateStore store, RestoreFunc restoreFunc) { restoreConsumer.unsubscribe(storePartition); } + public StateStore getStore(String name) { + return stores.get(name); + } + public void cleanup() throws IOException { // clean up any unknown files in the state directory for (File file : this.baseDir.listFiles()) { 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 0daa054657130..14bec7ea5fc65 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 @@ -24,7 +24,7 @@ import java.util.ArrayList; import java.util.List; -public class SinkNode extends ProcessorNode { +public class SinkNode extends ProcessorNode { private final String topic; private final Serializer keySerializer; diff --git a/stream/src/main/java/org/apache/kafka/streaming/processor/internals/SourceNode.java b/stream/src/main/java/org/apache/kafka/streaming/processor/internals/SourceNode.java index e86fe9c4d661d..f20afee0d3bef 100644 --- a/stream/src/main/java/org/apache/kafka/streaming/processor/internals/SourceNode.java +++ b/stream/src/main/java/org/apache/kafka/streaming/processor/internals/SourceNode.java @@ -20,11 +20,13 @@ import org.apache.kafka.common.serialization.Deserializer; import org.apache.kafka.streaming.processor.ProcessorContext; -public class SourceNode extends ProcessorNode { +public class SourceNode extends ProcessorNode { public Deserializer keyDeserializer; public Deserializer valDeserializer; + private ProcessorContext context; + public SourceNode(String name, Deserializer keyDeserializer, Deserializer valDeserializer) { super(name); @@ -34,15 +36,12 @@ public SourceNode(String name, Deserializer keyDeserializer, Deserializer @Override public void init(ProcessorContext context) { - // do nothing + this.context = context; } @Override public void process(K key, V value) { - // just forward to all children - for (ProcessorNode childNode : this.children()) { - childNode.process(key, value); - } + context.forward(key, value); } @Override @@ -50,4 +49,4 @@ public void close() { // do nothing } -} \ No newline at end of file +} diff --git a/stream/src/test/java/org/apache/kafka/streaming/KStreamWindowedTest.java b/stream/src/test/java/org/apache/kafka/streaming/KStreamWindowedTest.java index 07fccfdbd241b..85fcbec8d1825 100644 --- a/stream/src/test/java/org/apache/kafka/streaming/KStreamWindowedTest.java +++ b/stream/src/test/java/org/apache/kafka/streaming/KStreamWindowedTest.java @@ -21,7 +21,7 @@ import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.streaming.kstream.KStream; import org.apache.kafka.streaming.kstream.KStreamBuilder; -import org.apache.kafka.streaming.kstream.Window; +import org.apache.kafka.streaming.kstream.WindowDef; import org.apache.kafka.streaming.kstream.internals.KStreamSource; import org.apache.kafka.test.MockKStreamBuilder; import org.apache.kafka.test.MockProcessorContext; @@ -46,7 +46,7 @@ public void testWindowedStream() { final int[] expectedKeys = new int[]{0, 1, 2, 3}; KStream stream; - Window window; + WindowDef window; window = new UnlimitedWindow<>(); stream = topology.from(keyDeserializer, valDeserializer, topicName); diff --git a/stream/src/test/java/org/apache/kafka/streaming/internals/KStreamWindowedTest.java b/stream/src/test/java/org/apache/kafka/streaming/internals/KStreamWindowedTest.java index 88e92a551a981..aaa4e3a04a1f0 100644 --- a/stream/src/test/java/org/apache/kafka/streaming/internals/KStreamWindowedTest.java +++ b/stream/src/test/java/org/apache/kafka/streaming/internals/KStreamWindowedTest.java @@ -21,7 +21,7 @@ import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.streaming.kstream.KStream; import org.apache.kafka.streaming.kstream.KStreamBuilder; -import org.apache.kafka.streaming.kstream.Window; +import org.apache.kafka.streaming.kstream.WindowDef; import org.apache.kafka.streaming.kstream.internals.KStreamSource; import org.apache.kafka.test.MockKStreamBuilder; import org.apache.kafka.test.MockProcessorContext; @@ -46,7 +46,7 @@ public void testWindowedStream() { final int[] expectedKeys = new int[]{0, 1, 2, 3}; KStream stream; - Window window; + WindowDef window; window = new UnlimitedWindow<>(); stream = topology.from(keyDeserializer, valDeserializer, topicName); diff --git a/stream/src/test/java/org/apache/kafka/test/UnlimitedWindow.java b/stream/src/test/java/org/apache/kafka/test/UnlimitedWindow.java index bea636167e11a..2dc4c6cbed287 100644 --- a/stream/src/test/java/org/apache/kafka/test/UnlimitedWindow.java +++ b/stream/src/test/java/org/apache/kafka/test/UnlimitedWindow.java @@ -19,14 +19,14 @@ import org.apache.kafka.streaming.processor.ProcessorContext; import org.apache.kafka.streaming.kstream.KeyValue; -import org.apache.kafka.streaming.kstream.Window; +import org.apache.kafka.streaming.kstream.WindowDef; import org.apache.kafka.streaming.kstream.internals.FilteredIterator; import org.apache.kafka.streaming.processor.internals.Stamped; import java.util.Iterator; import java.util.LinkedList; -public class UnlimitedWindow implements Window { +public class UnlimitedWindow implements WindowDef { private LinkedList>> list = new LinkedList<>();