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

如何在Spring WebFlux中不使用@Scheduled注解实现特定时间定时任务?

仅用WebFlux API实现每日特定时间执行的定时任务

现有一段基于WebFlux的间隔式定时任务代码,需求是将其改为每日固定时间(例如晚10点)执行,且不依赖org.springframework.scheduling包,仅使用WebFlux原生API实现。

原代码

@EventListener(ApplicationReadyEvent.class)
public void scheduleUpload() {
    Flux.interval(schedulerConfigurationProperties.getUpload().getPeriod())
            .publishOn(Schedulers.boundedElastic())
            .delaySubscription(schedulerConfigurationProperties.getUpload().getInitialDelay())
            .onBackpressureBuffer(1, tick -> logger.debug("Drop tick {} in buffer", tick))
            .onBackpressureDrop(tick -> logger.debug("Drop tick {} in backpressure", tick))
            .concatMap(this::processReportUploadTick)
            .onErrorResume(this::onError)
            .subscribe();
}


public Mono<Boolean> processReportUploadTick(final Long tick) {
    return schedulerTickProcessor.processScheduledReportUploadTask(tick)
            .repeatWhen(this::hasNextUploadTask)
            .reduce(Boolean::logicalAnd)
            .doOnNext(v -> logger.debug("After processScheduledReportUploadTask tick {}", tick))
            .doFinally(s -> logger.debug("Tick finished {}", tick));
}

实现方案

完全可以通过WebFlux的Reactor API实现,核心思路是计算首次执行的初始延迟,之后固定每日间隔执行,具体步骤如下:

1. 编写初始延迟计算方法

先计算从当前时间到下一次目标执行时间的间隔,处理当天已过目标时间的情况:

private Duration calculateDailyInitialDelay() {
    // 设定每日执行时间:晚10点
    LocalTime targetTime = LocalTime.of(22, 0);
    // 若需指定时区,用 LocalDateTime.now(ZoneId.of("Asia/Shanghai"))
    LocalDateTime now = LocalDateTime.now();
    LocalDateTime nextRunTime = now.with(targetTime);

    // 如果当前时间已经过了当日目标时间,顺延到次日
    if (now.isAfter(nextRunTime)) {
        nextRunTime = nextRunTime.plusDays(1);
    }

    return Duration.between(now, nextRunTime);
}

2. 修改调度逻辑

使用Flux.interval的重载方法,传入初始延迟和每日间隔(24小时),替换原有的间隔+延迟逻辑:

@EventListener(ApplicationReadyEvent.class)
public void scheduleUpload() {
    Duration initialDelay = calculateDailyInitialDelay();
    // 每日间隔24小时
    Duration dailyInterval = Duration.ofHours(24);

    Flux.interval(initialDelay, dailyInterval)
            .publishOn(Schedulers.boundedElastic())
            .onBackpressureBuffer(1, tick -> logger.debug("Drop tick {} in buffer", tick))
            .onBackpressureDrop(tick -> logger.debug("Drop tick {} in backpressure", tick))
            .concatMap(this::processReportUploadTick)
            .onErrorResume(this::onError)
            .subscribe();
}

关键说明

  • 完全基于Reactor的Flux.intervalAPI,无需依赖spring-scheduling包
  • 初始延迟计算自动处理当日已过目标时间的场景,保证首次执行在最近的目标时间点
  • 原有背压处理、错误恢复、任务执行逻辑完全保留,无需修改processReportUploadTick方法
  • 若需要时区支持,在计算当前时间时指定对应时区即可,避免跨时区时间偏差

内容的提问来源于stack exchange,提问作者jerdys

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 06:05:19