【问题标题】:Spring cloud Kafka Stream StreamsUncaughtExceptionHandlerSpring Cloud Kafka Stream StreamsUncaughtExceptionHandler
【发布时间】:2022-03-02 08:46:05
【问题描述】:

我正在尝试将 StreamsUncaughtExceptionHandler 添加到我的 Kafka 流处理器中。这个处理器是用 Kafka 函数编写的。我查看了 suggestion provided by Artem Bilan 以将 StreamsUncaughtExceptionHandler 包含到我的服务中,但我的异常从未被它捕获/处理。

配置 Bean:

@Autowired
UnCaughtExceptionHandler exceptionHandler;

@Bean
public StreamsBuilderFactoryBeanConfigurer streamsCustomizer() {
    return new StreamsBuilderFactoryBeanConfigurer() {

        @Override
        public void configure(StreamsBuilderFactoryBean factoryBean) {;
            factoryBean.setStreamsUncaughtExceptionHandler(exceptionHandler);
        }

        @Override
        public int getOrder() {
            return Integer.MAX_VALUE;
        }

    };
}

自定义异常处理程序:

    @Component
public class UnCaughtExceptionHandler implements StreamsUncaughtExceptionHandler {

  @Autowired
  private StreamBridge streamBridge;

  @Override
  public StreamThreadExceptionResponse handle(Throwable exception) {
    return StreamThreadExceptionResponse.REPLACE_THREAD;
  }
}

流处理函数:

@Autowired
private MyService service;

@Bean
public Function<KStream<String, Input>, KStream<String, Output>> processor() {
    final AtomicReference<KeyValue<String, Output>> result = new AtomicReference<>(null);
    return kStream -> kStream
            .filter((key, value) -> value != null)
            .filter((key, value) -> {
                Optional<Output> outputResult = service.process(value);
                if (outputResult.isPresent()) {
                    result.set(new KeyValue<>(key, outputResult.get()));
                    return true;
                }
                return false;
            })
        .map((messageKey, messageValue) -> result.get());
}

我希望 UnCaughtExceptionHandler 能够处理 service.process() 方法抛出的任何异常。但是异常永远不会进入handle方法;相反,它们传播到根并死去客户端。也看过this solution,但我想以更独立的方式处理它。

问题:如何使用StreamsUncaughtExceptionHandler 处理任何处理异常?

  • Spring Boot 版本:2.6.3
  • spring-cloud-stream 版本:3.2.1
  • spring-cloud-stream-binder-kafka-streams: 3.2.1
  • kafka 流:3.0.0

可重现的例子:spring-cloud-kafka-streams-exception

【问题讨论】:

    标签: java apache-kafka apache-kafka-streams spring-kafka spring-cloud-stream


    【解决方案1】:

    您可以尝试以下几件事。

    1. 尝试在StreamsBuilderFactoryBean中的this line处设置断点,看看configure的值是多少。这应该会提供一些线索。

    2. 我注意到您在配置的 impl 中为订单设置了 Integer.MAX_VALUE。默认情况下,StreamsBuilderFactoryBean 使用Integer.MAX_VALUE - 1000 的阶段值,因此当工厂bean 准备好启动时,配置器可能还不可用,因为Integer.MAX_VALUE 的优先级较低。您可以将您的订单更改为 Integer.MAX_VALUE - 5000 之类的内容,以确保在启动工厂 bean 之前完全实例化配置 bean。

    从这些选项开始,看看它们是否能说明问题。如果仍然存在,请随时与我们分享一个可重现的小型示例应用程序。

    【讨论】:

    • @sobychako,在该行添加了一个断点并进行了检查。该值不为空。但是该类没有字段。按照建议更改了覆盖的 getOrder 配置值,但问题仍然存在。我将快速创建一个可重现的示例应用程序并在这里分享。
    • @sobychako,用可重现的示例应用程序更新了问题。
    • 感谢您的示例,这原来是活页夹中的错误。我们将不得不进行修复并在下一个版本中提供。在那之前你有机会使用活页夹的快照版本吗?
    • 感谢您的关注。是的,我们暂时可以使用SNAPSHOT版本。请在 SNAPSHOT 版本可用时添加评论。确认一下,是“spring-cloud-stream-binder-kafka”吧?
    • 不,是spring-cloud-stream-binder-kafka-streams。您还需要来自spring-kafka 的快照。我一定会在这里更新。
    猜你喜欢
    • 2018-04-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-15
    • 1970-01-01
    • 2019-02-05
    • 1970-01-01
    相关资源
    最近更新 更多