X-Git-Url: https://juplo.de/gitweb/?a=blobdiff_plain;f=src%2Fmain%2Fjava%2Fde%2Fjuplo%2Fkafka%2Fwordcount%2Ftop10%2FTop10StreamProcessor.java;h=343ab4d85e970199981798a89efea7714c70023c;hb=9e0ee342b531e7cfcddfe2e34390807802725a14;hp=2ff078c181225a87251a50bdfcd8f5dceb8a355e;hpb=b69e1270a308f200b2640a01d37d4636a0a549e1;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 2ff078c..343ab4d 100644 --- a/src/main/java/de/juplo/kafka/wordcount/top10/Top10StreamProcessor.java +++ b/src/main/java/de/juplo/kafka/wordcount/top10/Top10StreamProcessor.java @@ -1,11 +1,10 @@ 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; @@ -20,12 +19,13 @@ public class Top10StreamProcessor public Top10StreamProcessor( String inputTopic, String outputTopic, - Properties props) + Properties props, + KeyValueBytesStoreSupplier storeSupplier) { Topology topology = Top10StreamProcessor.buildTopology( inputTopic, outputTopic, - null); + storeSupplier); streams = new KafkaStreams(topology, props); } @@ -43,7 +43,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); @@ -55,7 +56,7 @@ public class Top10StreamProcessor ReadOnlyKeyValueStore getStore(String name) { - return null; + return streams.store(StoreQueryParameters.fromNameAndType(name, QueryableStoreTypes.keyValueStore())); } public void start()