集成Spring Boot和Reactor-Kafka的KafkaReceiver

ric*_*din 6 java spring spring-boot project-reactor spring-kafka

我正在尝试使用该库开发 Spring Boot 应用程序,reactor-kafka以对从 Kafka 主题读取的某些消息做出反应。

我有一个配置类,可以构建一个KafkaReceiver.

@Configuration
public class MyConfiguration {

    @Bean
    public KafkaReceiver<String, String> kafkaReceiver() {
        Map<String, Object> props = new HashMap<>();
        // Options initialisation...
        final ReceiverOptions<String, String> receiverOptions =
                ReceiverOptions.<String, string>create(props)
                               .subscription(Collections.singleton(consumer.getTopic()));
        return KafkaReceiver.create(receiverOptions);
    } 
}
Run Code Online (Sandbox Code Playgroud)

嗯……现在呢?使用不那么反应性的spring-kafka库,我可以注释一个方法,@KafkaListenerSpring Boot 将为我创建一个从 Kafka 主题监听的线程。

我应该把KafkaReceiver,放在哪里?在所有的例子中我发现直接使用该main方法,但这不是Boot way

我正在使用 Spring Boot 2.1.3 和 Reactor-Kafka 1.1.0

提前致谢。

Art*_*lan 7

既然你有了那个KafkaReceiverbean,现在你可以这样做:

@Bean
public ApplicationRunner runner(KafkaReceiver<String, String> kafkaReceiver) {
        return args -> {
                kafkaReceiver.receive()
                          ...
                          .sunbscribe();
        };
}
Run Code Online (Sandbox Code Playgroud)

当这颗ApplicationRunner豆子准备好时,就会被踢掉ApplicationContext。有关详细信息,请参阅其 JavaDocs。

  • 非常感谢。也许,您可以在文档中的某个位置添加此集成:) (2认同)