【发布时间】: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
【问题讨论】:
标签: java apache-kafka apache-kafka-streams spring-kafka spring-cloud-stream