【发布时间】:2021-04-13 18:56:24
【问题描述】:
这是我的应用程序,它使用来自 Kafka 主题的数据,然后将计算结果发送到一个主题。
@SpringBootApplication
@EnableBinding(KStreamProcessor.class)
public class WordCountProcessorApplication {
@StreamListener("input")
@SendTo("output")
public KStream<?, WordCount> process(KStream<?, String> input) {
return input
.flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+")))
.groupBy((key, value) -> value)
.windowedBy(TimeWindows.of(5000))
.count(Materialized.as("WordCounts-multi"))
.toStream()
.map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, new Date(key.window().start()), new Date(key.window().end()))));
}
public static void main(String[] args) {
SpringApplication.run(WordCountProcessorApplication.class, args);
}
如何在每次使用来自 Kafka 主题的数据之前打印日志?
【问题讨论】:
标签: apache-kafka-streams spring-cloud-stream