Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 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 @@ -18,7 +18,7 @@
package org.apache.kafka.streaming.examples;

import org.apache.kafka.streaming.KafkaStreaming;
import org.apache.kafka.streaming.processor.KafkaProcessor;
import org.apache.kafka.streaming.processor.Processor;
import org.apache.kafka.streaming.processor.TopologyBuilder;
import org.apache.kafka.streaming.StreamingConfig;
import org.apache.kafka.streaming.processor.ProcessorContext;
Expand All @@ -29,7 +29,7 @@

public class SimpleProcessJob {

private static class MyProcessor extends KafkaProcessor<String, Integer, Object, Object> {
private static class MyProcessor extends Processor<String, Integer, Object, Object> {
private ProcessorContext context;

public MyProcessor(String name) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@
import org.apache.kafka.common.serialization.IntegerDeserializer;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.streaming.KafkaStreaming;
import org.apache.kafka.streaming.processor.KafkaProcessor;
import org.apache.kafka.streaming.processor.Processor;
import org.apache.kafka.streaming.processor.TopologyBuilder;
import org.apache.kafka.streaming.StreamingConfig;
import org.apache.kafka.streaming.processor.ProcessorContext;
Expand All @@ -33,7 +33,7 @@

public class StatefulProcessJob {

private static class MyProcessor extends KafkaProcessor<String, Integer, Object, Object> {
private static class MyProcessor extends Processor<String, Integer, Object, Object> {
private ProcessorContext context;
private KeyValueStore<String, Integer> kvStore;

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

import org.apache.kafka.common.serialization.Deserializer;
import org.apache.kafka.common.serialization.Serializer;
import org.apache.kafka.streaming.processor.KafkaProcessor;
import org.apache.kafka.streaming.processor.Processor;

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,9 @@
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.RecordCollector;
import org.apache.kafka.streaming.processor.RestoreFunc;
import org.apache.kafka.streaming.processor.internals.Stamped;

Expand Down Expand Up @@ -173,7 +174,7 @@ public void flush() {
IntegerSerializer intSerializer = new IntegerSerializer();
ByteArraySerializer byteArraySerializer = new ByteArraySerializer();

RecordCollector collector = context.recordCollector();
RecordCollector collector = ((ProcessorContextImpl)context).recordCollector();

for (Map.Entry<K, ValueList<V>> entry : map.entrySet()) {
ValueList<V> values = entry.getValue();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,11 @@
package org.apache.kafka.streaming.kstream.internals;

import org.apache.kafka.common.KafkaException;
import org.apache.kafka.streaming.processor.KafkaProcessor;
import org.apache.kafka.streaming.processor.Processor;
import org.apache.kafka.streaming.processor.ProcessorMetadata;
import org.apache.kafka.streaming.kstream.Predicate;

class KStreamBranch<K, V> extends KafkaProcessor<K, V, K, V> {
class KStreamBranch<K, V> extends Processor<K, V, K, V> {

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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,11 @@

package org.apache.kafka.streaming.kstream.internals;

import org.apache.kafka.streaming.processor.KafkaProcessor;
import org.apache.kafka.streaming.processor.Processor;
import org.apache.kafka.streaming.kstream.Predicate;
import org.apache.kafka.streaming.processor.ProcessorMetadata;

class KStreamFilter<K, V> extends KafkaProcessor<K, V> {
class KStreamFilter<K, V> extends Processor<K, V> {

private final PredicateOut<K, V> predicateOut;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,12 +17,12 @@

package org.apache.kafka.streaming.kstream.internals;

import org.apache.kafka.streaming.processor.KafkaProcessor;
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;

class KStreamFlatMap<K1, V1, K2, V2> extends KafkaProcessor<K1, V1> {
class KStreamFlatMap<K1, V1, K2, V2> extends Processor<K1, V1> {

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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,11 @@

package org.apache.kafka.streaming.kstream.internals;

import org.apache.kafka.streaming.processor.KafkaProcessor;
import org.apache.kafka.streaming.processor.Processor;
import org.apache.kafka.streaming.kstream.ValueMapper;
import org.apache.kafka.streaming.processor.ProcessorMetadata;

class KStreamFlatMapValues<K1, V1, V2> extends KafkaProcessor<K1, V1, K1, V2> {
class KStreamFlatMapValues<K1, V1, V2> extends Processor<K1, V1, K1, V2> {

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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,6 @@
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.Processor;
import org.apache.kafka.streaming.processor.ProcessorMetadata;
import org.apache.kafka.streaming.processor.TopologyBuilder;
import org.apache.kafka.streaming.kstream.KStreamWindowed;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,14 +17,14 @@

package org.apache.kafka.streaming.kstream.internals;

import org.apache.kafka.streaming.processor.KafkaProcessor;
import org.apache.kafka.streaming.processor.Processor;
import org.apache.kafka.streaming.processor.ProcessorContext;
import org.apache.kafka.streaming.kstream.ValueJoiner;
import org.apache.kafka.streaming.kstream.Window;

import java.util.Iterator;

class KStreamJoin<K, V, V1, V2> extends KafkaProcessor<K, V1, K, V> {
class KStreamJoin<K, V, V1, V2> extends Processor<K, V1, K, V> {

private static final String JOIN_NAME = "KAFKA-JOIN";
private static final String JOIN_OTHER_NAME = "KAFKA-JOIN-OTHER";
Expand All @@ -38,7 +38,7 @@ private static abstract class Finder<K, T> {
private final Finder<K, V1> finder1;
private final Finder<K, V2> finder2;
private final ValueJoiner<V1, V2, V> joiner;
final KafkaProcessor<K, V2, K, V> processorForOtherStream;
final Processor<K, V2, K, V> processorForOtherStream;

private ProcessorContext context;

Expand Down Expand Up @@ -100,8 +100,8 @@ public void process(K key, V1 value) {
}
}

private KafkaProcessor<K, V2, K, V> processorForOther() {
return new KafkaProcessor<K, V2, K, V>(JOIN_OTHER_NAME) {
private Processor<K, V2, K, V> processorForOther() {
return new Processor<K, V2, K, V>(JOIN_OTHER_NAME) {

@SuppressWarnings("unchecked")
@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,13 +17,12 @@

package org.apache.kafka.streaming.kstream.internals;

import org.apache.kafka.streaming.processor.KafkaProcessor;
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.ProcessorContext;
import org.apache.kafka.streaming.processor.ProcessorMetadata;

class KStreamMap<K1, V1, K2, V2> extends KafkaProcessor<K1, V1> {
class KStreamMap<K1, V1, K2, V2> extends Processor<K1, V1> {

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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,11 @@

package org.apache.kafka.streaming.kstream.internals;

import org.apache.kafka.streaming.processor.KafkaProcessor;
import org.apache.kafka.streaming.processor.Processor;
import org.apache.kafka.streaming.kstream.ValueMapper;
import org.apache.kafka.streaming.processor.ProcessorMetadata;

class KStreamMapValues<K1, V1, V2> extends KafkaProcessor<K1, V1, K1, V2> {
class KStreamMapValues<K1, V1, V2> extends Processor<K1, V1, K1, V2> {

private final ValueMapper<V1, V2> mapper;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,10 @@

package org.apache.kafka.streaming.kstream.internals;

import org.apache.kafka.streaming.processor.KafkaProcessor;
import org.apache.kafka.streaming.processor.Processor;
import org.apache.kafka.streaming.processor.ProcessorMetadata;

public class KStreamProcessor<K, V> extends KafkaProcessor<K, V> {
public class KStreamProcessor<K, V> extends Processor<K, V> {

private final Processor<K, V> processor;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,11 @@
package org.apache.kafka.streaming.kstream.internals;

import org.apache.kafka.common.serialization.Serializer;
import org.apache.kafka.streaming.processor.KafkaProcessor;
import org.apache.kafka.streaming.processor.Processor;
import org.apache.kafka.streaming.processor.ProcessorMetadata;
import org.apache.kafka.streaming.processor.ProcessorContext;

class KStreamSend<K, V> extends KafkaProcessor<K, V, Object, Object> {
class KStreamSend<K, V> extends Processor<K, V, Object, Object> {

private ProcessorContext context;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,15 +17,15 @@

package org.apache.kafka.streaming.kstream.internals;

import org.apache.kafka.streaming.processor.KafkaProcessor;
import org.apache.kafka.streaming.processor.Processor;
import org.apache.kafka.streaming.processor.TopologyBuilder;
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<K, V> extends KafkaProcessor<K, V, K, V> {
public class KStreamWindow<K, V> extends Processor<K, V, K, V> {

public static final class KStreamWindowedImpl<K, V> extends KStreamImpl<K, V> implements KStreamWindowed<K, V> {

Expand Down

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -23,5 +23,7 @@ public interface Processor<K, V> {

void process(K key, V value);

void punctuate(long streamTime);

void close();
}
Original file line number Diff line number Diff line change
Expand Up @@ -63,13 +63,6 @@ public interface ProcessorContext {
*/
Deserializer<?> valueDeserializer();

/**
* Returns a RecordCollector
*
* @return RecordCollector
*/
RecordCollector recordCollector();

/**
* Returns the state directory for the partition.
*
Expand All @@ -91,19 +84,10 @@ public interface ProcessorContext {
*/
void register(StateStore store, RestoreFunc restoreFunc);

/**
* Flush the local state of this context
*/
void flush();
void schedule(Processor processor, long interval);

<K, V> void forward(K key, V value);

void send(String topic, Object key, Object value);

void send(String topic, Object key, Object value, Serializer<Object> keySerializer, Serializer<Object> valSerializer);

void schedule(KafkaProcessor processor, long interval);

void commit();

String topic();
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 Punctuator {
public interface ProcessorFactory {

void punctuate(long streamTime);
Processor build();
}

This file was deleted.

Loading