如何在Reactor Scheduler中实现元素发射延迟并保证接口调用间隔≥1秒
解决方案
现有代码的问题在于delaySubscription仅延迟了订阅外部方法的时机,而非确保两个外部API调用的执行间隔至少1秒。如果method1和method2被近乎同时调用,两个外部方法仍会近乎同时执行,无法满足间隔要求。
要实现两个外部API调用的处理间隔至少1秒,我们可以通过维护一个全局的"上次调用完成信号",让每次新的调用都等待上一次调用完成后再延迟1秒执行。具体实现如下:
@Service @RequiredArgsConstructor public class ServiceFoo { private final ServiceBar serviceBar; // 用AtomicReference保存上一次调用完成的Mono,初始为已完成的空信号 private final AtomicReference<Mono<Void>> lastCallSignal = new AtomicReference<>(Mono.empty()); private final Scheduler scheduler = Schedulers.newSingle("external-api-calls"); public Mono<String> method1(String string) { return executeWithInterval(() -> serviceBar.externalMethod1(string)); } public Mono<String> method2(String string) { return executeWithInterval(() -> serviceBar.externalMethod2(string)); } private Mono<String> executeWithInterval(Supplier<Mono<String>> apiCallSupplier) { // 原子更新调用信号:先等上次调用完成+1秒延迟,再执行当前API调用 Mono<Void> newSignal = lastCallSignal.updateAndGet(prev -> prev.then(Mono.delay(Duration.ofSeconds(1), scheduler)) .then(apiCallSupplier.get()) .then() ); // 返回当前API调用的结果 return newSignal.then(apiCallSupplier.get()); } }
核心逻辑说明
AtomicReference<Mono<Void>> lastCallSignal:线程安全地跟踪上一次外部API调用完成的信号,避免并发调用时的冲突。executeWithInterval统一处理间隔逻辑:- 等待上一次调用完成后,延迟1秒
- 执行当前的外部API调用
- 更新全局信号为当前调用完成的标记
- 无论
method1和method2的调用顺序、间隔如何,都能保证两个外部API的执行间隔至少为1秒。
内容的提问来源于stack exchange,提问作者Dmirii Berrington
相关产品推荐
相关产品推荐

