如何用Reactor实现定时执行并响应外部信号的任务调度?
解决Reactor中“定时+外部触发”的任务调度问题
核心思路是构建一个串行循环:每次循环中要么等待定时时间到达,要么等待外部触发信号,执行任务后进入下一轮循环。同时通过热Flux的重放特性避免信号丢失,通过串行执行避免任务排队堆积。
实现步骤与代码示例
- 创建外部触发的热信号源
使用replay().latest()保证新订阅总能拿到最新的触发信号,避免循环重新订阅时丢失信号:
import reactor.core.publisher.Sinks; // 外部触发信号源:重放最新信号,避免订阅间隙丢失 Sinks.Many<Void> triggerSink = Sinks.many().replay().latest();
- 初始化定时调度基准
用原子变量记录下一次定时执行的时间戳,确保外部触发后不重置计时器:
import java.util.concurrent.atomic.AtomicLong; AtomicLong nextScheduledTime = new AtomicLong(System.currentTimeMillis() + 10000);
- 定义任务逻辑
替换成你的实际任务代码:
Runnable task = () -> { // 此处编写你的任务逻辑 System.out.println("任务执行于:" + System.currentTimeMillis()); };
- 构建调度循环
通过Flux.defer()+repeat()实现串行循环,每次循环等待「定时到期」或「外部触发」任一事件:
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.time.Duration; Flux.defer(() -> { long now = System.currentTimeMillis(); // 计算距离下一次定时的剩余时间(可能为0,即已到定时时间) long remainingDelay = Math.max(0, nextScheduledTime.get() - now); // 定时信号:等待剩余时间到期 Mono<Void> scheduledSignal = Mono.delay(Duration.ofMillis(remainingDelay)).then(); // 外部触发信号:获取最新的触发信号 Mono<Void> triggerSignal = triggerSink.asFlux().next(); // 等待任一信号触发 return Mono.firstWithValue(scheduledSignal, triggerSignal) .doOnSuccess(__ -> { // 执行任务 task.run(); // 仅当本次是定时触发时,更新下一次定时时间 if (System.currentTimeMillis() >= nextScheduledTime.get()) { nextScheduledTime.addAndGet(10000); } }); }) .repeat() // 任务执行完后进入下一轮循环 .subscribe();
关键特性说明
- 无信号堆积:循环串行执行,任务执行时的外部触发信号仅保留最新一个,不会形成排队
- 不重置计时器:仅在定时触发时更新下一次定时时间,外部触发后仍沿用原定时周期
- 无信号丢失:
replay().latest()确保循环重新订阅时能拿到订阅间隙产生的最新触发信号 - 适配任务耗时过长:若任务执行耗时超过10秒,下一轮循环会立即执行(因为剩余延迟为0),保证两次执行间隔不超过10秒
外部触发调用
需要触发任务时,调用triggerSink.tryEmitEmpty()即可:
// 外部触发任务立即执行 triggerSink.tryEmitEmpty();
内容的提问来源于stack exchange,提问作者dnault
相关产品推荐
相关产品推荐

