You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 14:40:16