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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -133,7 +133,7 @@ public interface KStream<K, V> {
/**
* Processes all elements in this stream by applying a processor.
*
* @param processorFactory the class of ProcessorFactory
* @param processorDef the class of ProcessorDef
*/
<K1, V1> KStream<K1, V1> process(ProcessorFactory processorFactory);
<K1, V1> KStream<K1, V1> process(ProcessorDef processorDef);
}
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,7 @@ public String name() {
}

@Override
public Window<K, V> build() {
public Window<K, V> instance() {
return new SlidingWindow();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,4 +33,4 @@ public interface Window<K, V> extends StateStore {
Iterator<V> findBefore(K key, long timestamp);

void put(K key, V value, long timestamp);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -21,5 +21,5 @@ public interface WindowDef<K, V> {

String name();

Window<K, V> build();
}
Window<K, V> instance();
}
Original file line number Diff line number Diff line change
Expand Up @@ -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<K, V> implements ProcessorFactory {
class KStreamBranch<K, V> implements ProcessorDef {

private final Predicate<K, V>[] predicates;

Expand All @@ -31,7 +31,7 @@ public KStreamBranch(Predicate... predicates) {
}

@Override
public Processor build() {
public Processor instance() {
return new KStreamBranchProcessor();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<K, V> implements ProcessorFactory {
class KStreamFilter<K, V> implements ProcessorDef {

private final Predicate<K, V> predicate;
private final boolean filterOut;
Expand All @@ -32,7 +32,7 @@ public KStreamFilter(Predicate<K, V> predicate, boolean filterOut) {
}

@Override
public Processor build() {
public Processor instance() {
return new KStreamFilterProcessor();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<K1, V1, K2, V2> implements ProcessorFactory {
class KStreamFlatMap<K1, V1, K2, V2> implements ProcessorDef {

private final KeyValueFlatMap<K1, V1, K2, V2> mapper;

Expand All @@ -31,7 +31,7 @@ class KStreamFlatMap<K1, V1, K2, V2> implements ProcessorFactory {
}

@Override
public Processor build() {
public Processor instance() {
return new KStreamFlatMapProcessor();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<K1, V1, V2> implements ProcessorFactory {
class KStreamFlatMapValues<K1, V1, V2> implements ProcessorDef {

private final ValueMapper<V1, ? extends Iterable<V2>> mapper;

Expand All @@ -31,7 +31,7 @@ class KStreamFlatMapValues<K1, V1, V2> implements ProcessorFactory {
}

@Override
public Processor build() {
public Processor instance() {
return new KStreamFlatMapValuesProcessor();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,9 @@

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.KeyValueFlatMap;
import org.apache.kafka.streaming.kstream.KStreamWindowed;
import org.apache.kafka.streaming.kstream.KeyValueMapper;
import org.apache.kafka.streaming.kstream.Predicate;
Expand Down Expand Up @@ -130,12 +130,12 @@ public <V1> KStream<K, V1> flatMapValues(ValueMapper<V, ? extends Iterable<V1>>
}

@Override
public KStreamWindowed<K, V> with(WindowDef<K, V> windowDef) {
public KStreamWindowed<K, V> with(WindowDef<K, V> window) {
String name = WINDOWED_NAME + INDEX.getAndIncrement();

topology.addProcessor(name, new KStreamWindow<>(windowDef), this.name);
topology.addProcessor(name, new KStreamWindow<>(window), this.name);

return new KStreamWindowedImpl<>(topology, name, windowDef);
return new KStreamWindowedImpl<>(topology, name, window);
}

@Override
Expand Down Expand Up @@ -185,10 +185,10 @@ public void sendTo(String topic, Serializer<K> keySerializer, Serializer<V> valS

@SuppressWarnings("unchecked")
@Override
public <K1, V1> KStream<K1, V1> process(final ProcessorFactory processorFactory) {
public <K1, V1> KStream<K1, V1> 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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<K, V, V1, V2> implements ProcessorFactory {
class KStreamJoin<K, V, V1, V2> implements ProcessorDef {

private static abstract class Finder<K, T> {
abstract Iterator<T> find(K key, long timestamp);
Expand All @@ -37,9 +37,9 @@ private static abstract class Finder<K, T> {
private final boolean prior;

private Processor processorForOtherStream = null;
public final ProcessorFactory processorFactoryForOtherStream = new ProcessorFactory() {
public final ProcessorDef processorDefForOtherStream = new ProcessorDef() {
@Override
public Processor build() {
public Processor instance() {
return processorForOtherStream;
}
};
Expand All @@ -52,7 +52,7 @@ public Processor build() {
}

@Override
public Processor build() {
public Processor instance() {
return new KStreamJoinProcessor();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<K1, V1, K2, V2> implements ProcessorFactory {
class KStreamMap<K1, V1, K2, V2> implements ProcessorDef {

private final KeyValueMapper<K1, V1, K2, V2> mapper;

Expand All @@ -31,7 +31,7 @@ public KStreamMap(KeyValueMapper<K1, V1, K2, V2> mapper) {
}

@Override
public Processor build() {
public Processor instance() {
return new KStreamMapProcessor();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<K1, V1, V2> implements ProcessorFactory {
class KStreamMapValues<K1, V1, V2> implements ProcessorDef {

private final ValueMapper<V1, V2> mapper;

Expand All @@ -30,7 +30,7 @@ public KStreamMapValues(ValueMapper<V1, V2> mapper) {
}

@Override
public Processor build() {
public Processor instance() {
return new KStreamMapProcessor();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,12 +18,12 @@
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<K, V> implements ProcessorFactory {
class KStreamPassThrough<K, V> implements ProcessorDef {

@Override
public Processor build() {
public Processor instance() {
return new KStreamPassThroughProcessor();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<K, V> implements ProcessorFactory {
public class KStreamWindow<K, V> implements ProcessorDef {

private final WindowDef<K, V> windowDef;

Expand All @@ -36,7 +36,7 @@ public WindowDef<K, V> window() {
}

@Override
public Processor build() {
public Processor instance() {
return new KStreamWindowProcessor();
}

Expand All @@ -47,7 +47,7 @@ private class KStreamWindowProcessor extends KStreamProcessor<K, V> {
@Override
public void init(ProcessorContext context) {
super.init(context);
this.window = windowDef.build();
this.window = windowDef.instance();
this.window.init(context);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ private <V1, V2> KStream<K, V2> join(KStreamWindowed<K, V1> 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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@

package org.apache.kafka.streaming.processor;

public interface ProcessorFactory {
public interface ProcessorDef {

Processor build();
Processor instance();
}
Original file line number Diff line number Diff line change
Expand Up @@ -47,16 +47,16 @@ private interface NodeFactory {
private class ProcessorNodeFactory implements NodeFactory {
public final String[] parents;
private final String name;
private final ProcessorFactory factory;
private final ProcessorDef definition;

public ProcessorNodeFactory(String name, String[] parents, ProcessorFactory factory) {
public ProcessorNodeFactory(String name, String[] parents, ProcessorDef definition) {
this.name = name;
this.parents = parents.clone();
this.factory = factory;
this.definition = definition;
}

public ProcessorNode build() {
Processor processor = factory.build();
Processor processor = definition.instance();
return new ProcessorNode(name, processor);
}
}
Expand Down Expand Up @@ -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 definition, String... parentNames) {
if (nodeNames.contains(name))
throw new IllegalArgumentException("Processor " + name + " is already added.");

Expand All @@ -150,7 +150,7 @@ public final void addProcessor(String name, ProcessorFactory factory, String...
}

nodeNames.add(name);
nodeFactories.add(new ProcessorNodeFactory(name, parentNames, factory));
nodeFactories.add(new ProcessorNodeFactory(name, parentNames, definition));
}

/**
Expand Down Expand Up @@ -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 definition class: " + factory.getClass().getName());
}
}
} catch (Exception e) {
Expand Down
Loading