【问题标题】:How to configure an UncaughtExceptionHandler in Spring Kafka Stream如何在 Spring Kafka Stream 中配置 UncaughtExceptionHandler
【发布时间】: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


    【解决方案1】:

    是的。在那个旧版本中,您必须自己指定一个 KStreamBuilderFactoryBean bean,并进行适当的注入,而正是 KafkaStreamsDefaultConfiguration.DEFAULT_KSTREAM_BUILDER_BEAN_NAME

    在以后的版本中,我们已经有了一个 StreamsBuilderFactoryBeanConfigurer 来保持自动配置的 KStreamBuilderFactoryBean,但可以根据需要进行任何修改。

    更新

    您只需在应用程序上下文中将其创建为 bean,框架就会将其拾取并应用于 StreamsBuilderFactoryBean

        @Bean
        StreamsBuilderFactoryBeanConfigurer streamsCustomizer() {
            return new StreamsBuilderFactoryBeanConfigurer() {
    
                @Override
                public void configure(StreamsBuilderFactoryBean factoryBean) {
                    factoryBean.setCloseTimeout(...);
                }
    
                @Override
                public int getOrder() {
                    return Integer.MAX_VALUE;
                }
    
            };
        }
    

    【讨论】:

    • 我们正在从 1.5.15 版本迁移到 Spring Boot 2.2.6。现在,我们了解了使用接口StreamsBuilderFactoryBeanCustomizer 的解决方案的强大功能。感谢您的参考!
    • 您能否提供一个简短的例子来说明如何使用StreamsBuilderFactoryBeanCustomizer?谢谢
    • 在我的回答中查看更新。
    • @ArtemBilan 我正在使用带有功能编程的 spring cloud stream kafka。我尝试按照您上面建议的方式执行此操作。我确实将 StreamsUncaughtExceptionHandler 设置为 factoryBean 期望它会处理我的服务抛出的所有未捕获的异常。但它没有用。你知道函数式编程的任何工作示例吗? spring-cloud-stream-binder-kafka-streams:3.2.1 spring-cloud-stream:3.2.1 Spring-boot:2.6.3
    • 我建议您提出一个新的 SO 问题,提供更多详细信息。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-11-24
    • 2018-04-28
    • 2020-01-01
    相关资源
    最近更新 更多