X-Git-Url: https://juplo.de/gitweb/?a=blobdiff_plain;f=src%2Fmain%2Fjava%2Fde%2Fjuplo%2Fkafka%2Fwordcount%2Ftop10%2FTop10StreamProcessor.java;h=70ead8796c49617037fa39a176aa44390417a272;hb=3b7fae76b8abb62a8cae3a4a32c880b29bce0574;hp=d3846d89beab0cb573ddca8dcc4ffc477afee2fd;hpb=f350c2cafd9f0b290a021443cc7f5818974438e9;p=demos%2Fkafka%2Fwordcount diff --git a/src/main/java/de/juplo/kafka/wordcount/top10/Top10StreamProcessor.java b/src/main/java/de/juplo/kafka/wordcount/top10/Top10StreamProcessor.java index d3846d8..70ead87 100644 --- a/src/main/java/de/juplo/kafka/wordcount/top10/Top10StreamProcessor.java +++ b/src/main/java/de/juplo/kafka/wordcount/top10/Top10StreamProcessor.java @@ -1,10 +1,11 @@ package de.juplo.kafka.wordcount.top10; import lombok.extern.slf4j.Slf4j; -import org.apache.kafka.streams.KafkaStreams; -import org.apache.kafka.streams.KeyValue; -import org.apache.kafka.streams.StreamsBuilder; -import org.apache.kafka.streams.Topology; +import org.apache.kafka.streams.*; +import org.apache.kafka.streams.kstream.Materialized; +import org.apache.kafka.streams.state.KeyValueBytesStoreSupplier; +import org.apache.kafka.streams.state.QueryableStoreTypes; +import org.apache.kafka.streams.state.ReadOnlyKeyValueStore; import java.util.Properties; @@ -12,24 +13,29 @@ import java.util.Properties; @Slf4j public class Top10StreamProcessor { + public static final String STORE_NAME= "top10"; + public final KafkaStreams streams; public Top10StreamProcessor( String inputTopic, String outputTopic, - Properties props) + Properties props, + KeyValueBytesStoreSupplier storeSupplier) { Topology topology = Top10StreamProcessor.buildTopology( inputTopic, - outputTopic); + outputTopic, + storeSupplier); streams = new KafkaStreams(topology, props); } static Topology buildTopology( String inputTopic, - String outputTopic) + String outputTopic, + KeyValueBytesStoreSupplier storeSupplier) { StreamsBuilder builder = new StreamsBuilder(); @@ -39,7 +45,8 @@ public class Top10StreamProcessor .groupByKey() .aggregate( () -> new Ranking(), - (user, entry, ranking) -> ranking.add(entry)) + (user, entry, ranking) -> ranking.add(entry), + Materialized.as(storeSupplier)) .toStream() .to(outputTopic); @@ -49,6 +56,11 @@ public class Top10StreamProcessor return topology; } + ReadOnlyKeyValueStore getStore() + { + return streams.store(StoreQueryParameters.fromNameAndType(STORE_NAME, QueryableStoreTypes.keyValueStore())); + } + public void start() { log.info("Starting Stream-Processor");