Spring Webflux应用中Reactor Kafka消费流健康检查及Flux状态获取
如何查询Flux运行状态实现健康检查
Reactor没有提供内置的API直接查询Flux的运行状态,你可以通过自定义状态标记的方式实现状态感知,适配Spring Boot健康检查端点的需求:
- 首先定义状态枚举和原子状态变量,示例代码如下:
public enum ConsumeStatus { RUNNING, COMPLETED, ERROR } // 全局状态变量 private final AtomicReference<ConsumeStatus> consumeStatus = new AtomicReference<>(ConsumeStatus.RUNNING);
- 给Kafka消费的Flux挂载生命周期钩子,更新状态:
kafkaReceiver.receive() .doOnSubscribe(sub -> consumeStatus.set(ConsumeStatus.RUNNING)) .doOnComplete(() -> consumeStatus.set(ConsumeStatus.COMPLETED)) .doOnError(err -> consumeStatus.set(ConsumeStatus.ERROR)) // 你的正常消费逻辑 .subscribe();
- 自定义Spring Boot健康指示器,对外暴露状态:
@Component public class KafkaConsumerHealthIndicator implements ReactiveHealthIndicator { @Override public Mono<Health> health() { if (consumeStatus.get() == ConsumeStatus.RUNNING) { return Mono.just(Health.up().build()); } return Mono.just(Health.down().withDetail("status", consumeStatus.get()).build()); } }
配置完成后就能通过/actuator/health端点感知消费Flux的运行状态,下游监控/容器平台可以基于这个端点的返回结果执行重启操作。
为什么测试代码中doOnTerminate、doFinally未执行
你测试时抛出的OutOfMemoryError属于JVM致命错误(Fatal Error),Reactor默认的异常处理链路只会捕获Exception类型的非致命异常,碰到Error类型的错误会直接向上抛到调用线程,终止线程执行,因此走不到后续挂载的生命周期操作符。
如果要测试生命周期钩子的执行逻辑,要么改用普通的RuntimeException测试,要么显式指定捕获所有Throwable类型的错误,修正后的代码如下:
Flux.range(1, 5) .flatMap(record -> Mono.just(record) .map(i -> { throw new OutOfMemoryError("Forcing exception for " + i); }) .doOnNext(i -> System.out.println("doOnNext: " + i)) .doOnError(e -> System.err.println(e)) // 显式指定捕获所有Throwable类型 .onErrorResume(Throwable.class, e -> Mono.empty())) .doFinally(signalType -> System.err.println("doFinally: Terminating with Signal type: " + signalType)) .doOnTerminate(()-> System.err.println("doOnTerminate: executed")) // 订阅时指定全局错误消费者,避免错误直接抛到线程 .subscribe(null, e -> System.err.println("全局捕获错误: " + e));
注意:生产环境不建议捕获OutOfMemoryError这类致命错误,出现该错误时JVM已处于不稳定状态,直接终止进程重启是更合理的处理方式。
内容的提问来源于stack exchange,提问作者ankush
相关产品推荐
相关产品推荐

