From 5bac8ca829d6dc55eabfad4b788f798226e8b9bb Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Tue, 22 Sep 2015 15:06:42 -0700 Subject: [PATCH] comments addressed --- build.gradle | 2 +- .../org/apache/kafka/streams/kstream/KStream.java | 2 +- .../kstream/internals/KStreamFlatMapValues.java | 1 - .../streams/kstream/internals/KStreamImpl.java | 14 +++++++------- 4 files changed, 9 insertions(+), 10 deletions(-) diff --git a/build.gradle b/build.gradle index 3de9101e7cde6..b0e1b58c20dd6 100644 --- a/build.gradle +++ b/build.gradle @@ -538,7 +538,7 @@ project(':streams') { compile "$slf4jlog4j" compile 'org.rocksdb:rocksdbjni:3.10.1' - testCompile 'junit:junit:4.6' + testCompile '$junit' testCompile project(path: ':clients', configuration: 'archives') } diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/KStream.java b/streams/src/main/java/org/apache/kafka/streams/kstream/KStream.java index e266be6e5864f..7f101ab48af91 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/KStream.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/KStream.java @@ -53,7 +53,7 @@ public interface KStream { KStream map(KeyValueMapper> mapper); /** - * Creates a new stream by transforming valuesa by a mapper to all values of this stream + * Creates a new stream by transforming values by a mapper to all values of this stream * * @param mapper the instance of ValueMapper * @param the value type of the new stream diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamFlatMapValues.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamFlatMapValues.java index a56d1d625c34e..4f3e4755281b4 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamFlatMapValues.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamFlatMapValues.java @@ -25,7 +25,6 @@ class KStreamFlatMapValues implements ProcessorDef { private final ValueMapper> mapper; - @SuppressWarnings("unchecked") KStreamFlatMapValues(ValueMapper> mapper) { this.mapper = mapper; } diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamImpl.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamImpl.java index f972a551087dd..69366481cd7b2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamImpl.java @@ -52,6 +52,8 @@ public class KStreamImpl implements KStream { private static final String WINDOWED_NAME = "KAFKA-WINDOWED-"; + private static final String SINK_NAME = "KAFKA-SINK-"; + public static final String JOINTHIS_NAME = "KAFKA-JOINTHIS-"; public static final String JOINOTHER_NAME = "KAFKA-JOINOTHER-"; @@ -60,12 +62,10 @@ public class KStreamImpl implements KStream { public static final String SOURCE_NAME = "KAFKA-SOURCE-"; - public static final String SEND_NAME = "KAFKA-SEND-"; - public static final AtomicInteger INDEX = new AtomicInteger(1); - protected TopologyBuilder topology; - protected String name; + protected final TopologyBuilder topology; + protected final String name; public KStreamImpl(TopologyBuilder topology, String name) { this.topology = topology; @@ -160,7 +160,7 @@ public KStream through(String topic, Serializer valSerializer, Deserializer keyDeserializer, Deserializer valDeserializer) { - String sendName = SEND_NAME + INDEX.getAndIncrement(); + String sendName = SINK_NAME + INDEX.getAndIncrement(); topology.addSink(sendName, topic, keySerializer, valSerializer, this.name); @@ -178,14 +178,14 @@ public KStream through(String topic) { @Override public void to(String topic) { - String name = SEND_NAME + INDEX.getAndIncrement(); + String name = SINK_NAME + INDEX.getAndIncrement(); topology.addSink(name, topic, this.name); } @Override public void to(String topic, Serializer keySerializer, Serializer valSerializer) { - String name = SEND_NAME + INDEX.getAndIncrement(); + String name = SINK_NAME + INDEX.getAndIncrement(); topology.addSink(name, topic, keySerializer, valSerializer, this.name); }