mirror of https://github.com/apache/kafka.git
MINOR: Add missing imports to 'Hello Kafka Streams' examples (#4535)
Reviewers: Matthias J. Sax <matthias@confluent.io>, Guozhang Wang <wangguoz@gmail.com>
This commit is contained in:
parent
ac267dc5ce
commit
d8eddc6e16
|
@ -153,10 +153,12 @@
|
|||
<div class="code-example__snippet b-java-8 selected">
|
||||
<pre class="brush: java;">
|
||||
import org.apache.kafka.common.serialization.Serdes;
|
||||
import org.apache.kafka.common.utils.Bytes;
|
||||
import org.apache.kafka.streams.KafkaStreams;
|
||||
import org.apache.kafka.streams.StreamsBuilder;
|
||||
import org.apache.kafka.streams.StreamsConfig;
|
||||
import org.apache.kafka.streams.Topology;
|
||||
import org.apache.kafka.streams.kstream.KStream;
|
||||
import org.apache.kafka.streams.kstream.KTable;
|
||||
import org.apache.kafka.streams.kstream.Materialized;
|
||||
import org.apache.kafka.streams.kstream.Produced;
|
||||
import org.apache.kafka.streams.state.KeyValueStore;
|
||||
|
@ -192,14 +194,16 @@
|
|||
<div class="code-example__snippet b-java-7">
|
||||
<pre class="brush: java;">
|
||||
import org.apache.kafka.common.serialization.Serdes;
|
||||
import org.apache.kafka.common.utils.Bytes;
|
||||
import org.apache.kafka.streams.KafkaStreams;
|
||||
import org.apache.kafka.streams.StreamsBuilder;
|
||||
import org.apache.kafka.streams.StreamsConfig;
|
||||
import org.apache.kafka.streams.Topology;
|
||||
import org.apache.kafka.streams.kstream.KStream;
|
||||
import org.apache.kafka.streams.kstream.KTable;
|
||||
import org.apache.kafka.streams.kstream.ValueMapper;
|
||||
import org.apache.kafka.streams.kstream.KeyValueMapper;
|
||||
import org.apache.kafka.streams.kstream.Materialized;
|
||||
import org.apache.kafka.streams.kstream.Produced;
|
||||
import org.apache.kafka.streams.kstream.ValueMapper;
|
||||
import org.apache.kafka.streams.state.KeyValueStore;
|
||||
|
||||
import java.util.Arrays;
|
||||
|
@ -249,9 +253,10 @@
|
|||
import java.util.concurrent.TimeUnit
|
||||
|
||||
import org.apache.kafka.common.serialization._
|
||||
import org.apache.kafka.common.utils.Bytes
|
||||
import org.apache.kafka.streams._
|
||||
import org.apache.kafka.streams.kstream.{KeyValueMapper, Materialized, Produced, ValueMapper}
|
||||
import org.apache.kafka.streams.state.KeyValueStore;
|
||||
import org.apache.kafka.streams.kstream.{KStream, KTable, Materialized, Produced}
|
||||
import org.apache.kafka.streams.state.KeyValueStore
|
||||
|
||||
import scala.collection.JavaConverters.asJavaIterableConverter
|
||||
|
||||
|
@ -273,7 +278,7 @@
|
|||
.flatMapValues(textLine => textLine.toLowerCase.split("\\W+").toIterable.asJava)
|
||||
.groupBy((_, word) => word)
|
||||
.count(Materialized.as("counts-store").asInstanceOf[Materialized[String, Long, KeyValueStore[Bytes, Array[Byte]]]])
|
||||
wordCounts.toStream().to("WordsWithCountsTopic", Produced.with(Serdes.String(), Serdes.Long()))
|
||||
wordCounts.toStream().to("WordsWithCountsTopic", Produced.`with`(Serdes.String(), Serdes.Long()))
|
||||
|
||||
val streams: KafkaStreams = new KafkaStreams(builder.build(), config)
|
||||
streams.start()
|
||||
|
|
Loading…
Reference in New Issue