Flux.interval执行遇异常停止,如何实现异常后仍定时执行?
解决Flux.interval定时任务异常后停止调度的问题
原代码中,onErrorResume作用于整个主流,当customService.process()抛出异常时,会直接终止整个Flux.interval流,导致后续定时事件不再发射。要实现异常后仍持续执行定时任务,需将异常处理限定在每个单次任务的子流中,避免影响主定时流。
修改后的代码
Flux.interval(Duration.ofSeconds(30)) .doOnNext(duration -> log.info("Process has been started")) .flatMap(duration -> customService.process() .onErrorResume(throwable -> { log.error("Process execution failed", throwable); return Flux.empty(); }) ) .subscribe();
核心修改说明
- 将异常处理逻辑从主流的
onErrorResume移到flatMap内部的customService.process()子流中 - 单次任务的异常只会终止当前子流,不会中断
Flux.interval的定时发射逻辑 - 可在
onErrorResume中添加异常日志,便于问题排查
可选优化:添加重试逻辑
如果需要在任务失败后进行重试,可在子流中加入retry操作:
Flux.interval(Duration.ofSeconds(30)) .doOnNext(duration -> log.info("Process has been started")) .flatMap(duration -> customService.process() .retry(2) // 失败后最多重试2次 .onErrorResume(throwable -> { log.error("Process failed after retries", throwable); return Flux.empty(); }) ) .subscribe();
内容的提问来源于stack exchange,提问作者Alexey
相关产品推荐
相关产品推荐

