From 59e00643217a983d990257d80e7761769461727b Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Wed, 5 Aug 2015 17:17:28 -0700 Subject: [PATCH 1/2] SingleProcessorTopology implements Processor --- .../apache/kafka/stream/examples/MapKStreamJob.java | 2 +- .../kafka/stream/examples/PrintKStreamJob.java | 9 ++++++--- .../kafka/stream/examples/StatefulKStreamJob.java | 9 ++++++--- .../stream/topology/SingleProcessorTopology.java | 13 +++---------- 4 files changed, 16 insertions(+), 17 deletions(-) diff --git a/stream/src/main/java/org/apache/kafka/stream/examples/MapKStreamJob.java b/stream/src/main/java/org/apache/kafka/stream/examples/MapKStreamJob.java index 04aa1070fdc5f..fbdc0fee79d39 100644 --- a/stream/src/main/java/org/apache/kafka/stream/examples/MapKStreamJob.java +++ b/stream/src/main/java/org/apache/kafka/stream/examples/MapKStreamJob.java @@ -34,7 +34,7 @@ public class MapKStreamJob extends KStreamTopology { @Override public void topology() { - // With overriden de-serializer + // With overridden de-serializer KStream stream1 = from(new StringDeserializer(), new StringDeserializer(), "topic1"); stream1.map(new KeyValueMapper() { diff --git a/stream/src/main/java/org/apache/kafka/stream/examples/PrintKStreamJob.java b/stream/src/main/java/org/apache/kafka/stream/examples/PrintKStreamJob.java index 1de39d5aaad2a..9408b04db5737 100644 --- a/stream/src/main/java/org/apache/kafka/stream/examples/PrintKStreamJob.java +++ b/stream/src/main/java/org/apache/kafka/stream/examples/PrintKStreamJob.java @@ -17,7 +17,6 @@ package org.apache.kafka.stream.examples; -import org.apache.kafka.stream.topology.Processor; import org.apache.kafka.stream.KStreamContext; import org.apache.kafka.stream.KafkaStreaming; import org.apache.kafka.stream.StreamingConfig; @@ -25,10 +24,14 @@ import java.util.Properties; -public class PrintKStreamJob implements Processor { +public class PrintKStreamJob extends SingleProcessorTopology { private KStreamContext context; + public PrintKStreamJob(String... topics) { + super(topics); + } + @Override public void init(KStreamContext context) { this.context = context; @@ -55,7 +58,7 @@ public void close() { public static void main(String[] args) { KafkaStreaming streaming = new KafkaStreaming( - new SingleProcessorTopology(PrintKStreamJob.class, args), + new PrintKStreamJob(args), new StreamingConfig(new Properties()) ); streaming.run(); diff --git a/stream/src/main/java/org/apache/kafka/stream/examples/StatefulKStreamJob.java b/stream/src/main/java/org/apache/kafka/stream/examples/StatefulKStreamJob.java index ccfb911d15048..d8be0c94ba32e 100644 --- a/stream/src/main/java/org/apache/kafka/stream/examples/StatefulKStreamJob.java +++ b/stream/src/main/java/org/apache/kafka/stream/examples/StatefulKStreamJob.java @@ -17,7 +17,6 @@ package org.apache.kafka.stream.examples; -import org.apache.kafka.stream.topology.Processor; import org.apache.kafka.stream.KStreamContext; import org.apache.kafka.stream.KafkaStreaming; import org.apache.kafka.stream.StreamingConfig; @@ -29,11 +28,15 @@ import java.util.Properties; -public class StatefulKStreamJob implements Processor { +public class StatefulKStreamJob extends SingleProcessorTopology { private KStreamContext context; private KeyValueStore kvStore; + public StatefulKStreamJob(String... topics) { + super(topics); + } + @Override public void init(KStreamContext context) { this.context = context; @@ -71,7 +74,7 @@ public void close() { public static void main(String[] args) { KafkaStreaming streaming = new KafkaStreaming( - new SingleProcessorTopology(StatefulKStreamJob.class, args), + new StatefulKStreamJob(args), new StreamingConfig(new Properties()) ); streaming.run(); diff --git a/stream/src/main/java/org/apache/kafka/stream/topology/SingleProcessorTopology.java b/stream/src/main/java/org/apache/kafka/stream/topology/SingleProcessorTopology.java index 4febc6e82949f..98aadf7ae5086 100644 --- a/stream/src/main/java/org/apache/kafka/stream/topology/SingleProcessorTopology.java +++ b/stream/src/main/java/org/apache/kafka/stream/topology/SingleProcessorTopology.java @@ -19,24 +19,17 @@ import org.apache.kafka.common.utils.Utils; -public class SingleProcessorTopology extends KStreamTopology { +abstract public class SingleProcessorTopology extends KStreamTopology implements Processor { - private final Class processorClass; private final String[] topics; - public SingleProcessorTopology(Class processorClass, String... topics) { - this.processorClass = processorClass; + public SingleProcessorTopology(String... topics) { this.topics = topics; } @SuppressWarnings("unchecked") @Override public void topology() { - from(topics).process(newProcessor()); - } - - @SuppressWarnings("unchecked") - private Processor newProcessor() { - return (Processor) Utils.newInstance(processorClass); + from(topics).process((Processor)this); } } From e0ff77fe058c523116e6381c2d009adda55a515d Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Thu, 6 Aug 2015 14:30:48 -0700 Subject: [PATCH 2/2] Address Yasu's comments --- .../stream/examples/PrintKStreamJob.java | 49 +++++++------ .../stream/examples/StatefulKStreamJob.java | 70 ++++++++++--------- .../topology/SingleProcessorTopology.java | 35 ---------- 3 files changed, 63 insertions(+), 91 deletions(-) delete mode 100644 stream/src/main/java/org/apache/kafka/stream/topology/SingleProcessorTopology.java diff --git a/stream/src/main/java/org/apache/kafka/stream/examples/PrintKStreamJob.java b/stream/src/main/java/org/apache/kafka/stream/examples/PrintKStreamJob.java index 9408b04db5737..48ca24557810a 100644 --- a/stream/src/main/java/org/apache/kafka/stream/examples/PrintKStreamJob.java +++ b/stream/src/main/java/org/apache/kafka/stream/examples/PrintKStreamJob.java @@ -20,45 +20,48 @@ import org.apache.kafka.stream.KStreamContext; import org.apache.kafka.stream.KafkaStreaming; import org.apache.kafka.stream.StreamingConfig; -import org.apache.kafka.stream.topology.SingleProcessorTopology; +import org.apache.kafka.stream.topology.KStreamTopology; +import org.apache.kafka.stream.topology.Processor; import java.util.Properties; -public class PrintKStreamJob extends SingleProcessorTopology { +public class PrintKStreamJob extends KStreamTopology { - private KStreamContext context; + private class MyProcessor implements Processor { + private KStreamContext context; - public PrintKStreamJob(String... topics) { - super(topics); - } + @Override + public void init(KStreamContext context) { + this.context = context; + } - @Override - public void init(KStreamContext context) { - this.context = context; - } + @Override + public void process(K key, V value) { + System.out.println("[" + key + ", " + value + "]"); - @Override - public void process(K key, V value) { - System.out.println("[" + key + ", " + value + "]"); + context.commit(); - context.commit(); + context.send("topic", key, value); + } - context.send("topic", key, value); - } + @Override + public void punctuate(long streamTime) { + // do nothing + } - @Override - public void punctuate(long streamTime) { - // do nothing + @Override + public void close() { + // do nothing + } } + @SuppressWarnings("unchecked") @Override - public void close() { - // do nothing - } + public void topology() { from("topic").process(new MyProcessor()); } public static void main(String[] args) { KafkaStreaming streaming = new KafkaStreaming( - new PrintKStreamJob(args), + new PrintKStreamJob(), new StreamingConfig(new Properties()) ); streaming.run(); diff --git a/stream/src/main/java/org/apache/kafka/stream/examples/StatefulKStreamJob.java b/stream/src/main/java/org/apache/kafka/stream/examples/StatefulKStreamJob.java index d8be0c94ba32e..aa042b5d3a867 100644 --- a/stream/src/main/java/org/apache/kafka/stream/examples/StatefulKStreamJob.java +++ b/stream/src/main/java/org/apache/kafka/stream/examples/StatefulKStreamJob.java @@ -17,6 +17,7 @@ package org.apache.kafka.stream.examples; +import org.apache.kafka.stream.KStream; import org.apache.kafka.stream.KStreamContext; import org.apache.kafka.stream.KafkaStreaming; import org.apache.kafka.stream.StreamingConfig; @@ -24,57 +25,60 @@ import org.apache.kafka.stream.state.InMemoryKeyValueStore; import org.apache.kafka.stream.state.KeyValueIterator; import org.apache.kafka.stream.state.KeyValueStore; -import org.apache.kafka.stream.topology.SingleProcessorTopology; +import org.apache.kafka.stream.topology.KStreamTopology; +import org.apache.kafka.stream.topology.Processor; import java.util.Properties; -public class StatefulKStreamJob extends SingleProcessorTopology { +public class StatefulKStreamJob extends KStreamTopology { - private KStreamContext context; - private KeyValueStore kvStore; + private class MyProcessor implements Processor { + private KStreamContext context; + private KeyValueStore kvStore; - public StatefulKStreamJob(String... topics) { - super(topics); - } + @Override + public void init(KStreamContext context) { + this.context = context; + this.context.schedule(this, 1000); - @Override - public void init(KStreamContext context) { - this.context = context; - this.context.schedule(this, 1000); + this.kvStore = new InMemoryKeyValueStore<>("local-state", context); + } - this.kvStore = new InMemoryKeyValueStore<>("local-state", context); - } + @Override + public void process(String key, Integer value) { + Integer oldValue = this.kvStore.get(key); + if (oldValue == null) { + this.kvStore.put(key, value); + } else { + int newValue = oldValue + value; + this.kvStore.put(key, newValue); + } - @Override - public void process(String key, Integer value) { - Integer oldValue = this.kvStore.get(key); - if (oldValue == null) { - this.kvStore.put(key, value); - } else { - int newValue = oldValue + value; - this.kvStore.put(key, newValue); + context.commit(); } - context.commit(); - } + @Override + public void punctuate(long streamTime) { + KeyValueIterator iter = this.kvStore.all(); + while (iter.hasNext()) { + Entry entry = iter.next(); + System.out.println("[" + entry.key() + ", " + entry.value() + "]"); + } + } - @Override - public void punctuate(long streamTime) { - KeyValueIterator iter = this.kvStore.all(); - while (iter.hasNext()) { - Entry entry = iter.next(); - System.out.println("[" + entry.key() + ", " + entry.value() + "]"); + @Override + public void close() { + // do nothing } } + @SuppressWarnings("unchecked") @Override - public void close() { - // do nothing - } + public void topology() { ((KStream) from("topic")).process(new MyProcessor()); } public static void main(String[] args) { KafkaStreaming streaming = new KafkaStreaming( - new StatefulKStreamJob(args), + new StatefulKStreamJob(), new StreamingConfig(new Properties()) ); streaming.run(); diff --git a/stream/src/main/java/org/apache/kafka/stream/topology/SingleProcessorTopology.java b/stream/src/main/java/org/apache/kafka/stream/topology/SingleProcessorTopology.java deleted file mode 100644 index 98aadf7ae5086..0000000000000 --- a/stream/src/main/java/org/apache/kafka/stream/topology/SingleProcessorTopology.java +++ /dev/null @@ -1,35 +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.stream.topology; - -import org.apache.kafka.common.utils.Utils; - -abstract public class SingleProcessorTopology extends KStreamTopology implements Processor { - - private final String[] topics; - - public SingleProcessorTopology(String... topics) { - this.topics = topics; - } - - @SuppressWarnings("unchecked") - @Override - public void topology() { - from(topics).process((Processor)this); - } -}