【问题标题】:Spring Cloud Stream’s Apache KafkaSpring Cloud Stream 的 Apache Kafka
【发布时间】: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


    【解决方案1】:

    在拓扑的开头和结尾添加Transformer

    请参阅this discussion,其中有一个请求由框架自动将自定义转换器添加到拓扑中。

    决定添加您自己的解决方法就足够了。

    【讨论】:

    • ConsumerInterceptor 可以帮我实现我提到的场景吗?
    • 是的,这也可以(和ProducerInterceptor出站)。
    猜你喜欢
    • 2018-04-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-09-03
    • 2019-02-15
    • 1970-01-01
    相关资源
    最近更新 更多