【发布时间】:2020-05-04 14:27:35
【问题描述】:
我正在使用 spring-kafka 来实现一个使用 Spring Boot 1.5.16 的流应用程序。我们使用的spring-kafka 版本是1.3.8.RELEASE。
我正在寻找一种方法来关闭启动应用程序,以防出现终止与 Kafka Streams 关联的所有线程的错误。我发现在KafkaStreams 中可以注册未捕获异常的句柄。方法是setGlobalStateRestoreListener。
我看到这个方法暴露在spring-kafka 中,类型为KStreamBuilderFactoryBean。
我的问题如下。有没有一种简单的方法可以将UncaughtExceptionHandler 注册为 bean 并让 Spring 在工厂 bean 中正确注入?还是我应该自己创建KStreamBuilderFactoryBean 并手动设置处理程序?
@Bean(name = KafkaStreamsDefaultConfiguration.DEFAULT_KSTREAM_BUILDER_BEAN_NAME)
public KStreamBuilderFactoryBean kStreamBuilderFactoryBean(StreamsConfig streamsConfig) {
final KStreamBuilderFactoryBean streamBuilderFactoryBean = new KStreamBuilderFactoryBean(
streamsConfig);
streamBuilderFactoryBean.setUncaughtExceptionHandler((threadInError, exception) -> {
// Something happens here
});
return streamBuilderFactoryBean;
}
非常感谢。
【问题讨论】:
标签: java spring-batch apache-kafka-streams spring-kafka