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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 13:45:08