如何在Reactive Stream内实现每1分钟执行一次写入Kafka流程的调度
实现方案
你可以用Reactor内置的Flux.interval作为定时触发源,在流内部完成每分钟调度的逻辑,不需要依赖外部定时组件。首先需要先修正你原有代码的小问题:sendData方法声明了返回Flux<M>但内部直接调用了subscribe且没有返回值,我们调整为把流返回,交由上层调度逻辑订阅即可。
完整实现代码如下:
import reactor.core.publisher.Flux; import reactor.core.scheduler.Schedulers; import java.time.Duration; // 修正后的处理流方法,不再内部订阅,返回处理流供调度层调用 public Flux<M> sendDataFlow(){ return Flux.fromIterable(o.getC()) .publishOn(Schedulers.boundedElastic()) .flatMap(id -> Flux.fromIterable(getM(id)) .publishOn(Schedulers.boundedElastic()) .flatMap(n -> Flux.fromIterable(o.getD()) .publishOn(Schedulers.boundedElastic()) .flatMap(d -> Flux.just(sendToKafka)) ) ) .doOnError(throwable -> log.debug("Error while reading data : {} ", throwable.getMessage())) // 单次执行出错时兜底返回空,避免打断整个定时调度流 .onErrorResume(e -> Flux.empty()); } // 加入每分钟调度逻辑 public void run(String... args){ // 第一个参数为首次执行延迟,设为ZERO表示启动后立即执行第一次,之后每间隔1分钟触发一次 Flux.interval(Duration.ZERO, Duration.ofMinutes(1)) .flatMap(tick -> sendDataFlow()) .subscribe(); }
补充说明
- 如果不需要启动后立即执行第一次,把
Flux.interval的第一个参数调整为对应延迟时长即可,比如Duration.ofMinutes(1)就表示首次执行延迟1分钟,之后每次间隔1分钟。 - 错误兜底逻辑放在单次执行流内部,保证单次任务执行报错不会终止整个定时调度流,后续每分钟的任务仍能正常触发。
- 所有调度逻辑都在Reactor流内部完成,不需要引入Quartz、Spring
@Scheduled这类外部定时组件。
内容的提问来源于stack exchange,提问作者perplexedDev
相关产品推荐
相关产品推荐

