如何在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
相关产品推荐
相关产品推荐

