在对接第三方服务、异步任务调度等场景中,我们经常需要不断查询外部系统的状态,直到目标状态出现或者达到超时条件。使用Reactor Mono可以很优雅地实现这种异步轮询逻辑,避免传统阻塞式轮询带来的资源浪费问题。

核心实现思路
Reactor Mono实现异步轮询的核心逻辑是:先发起一次状态查询,判断是否满足终止条件,如果不满足就延迟一段时间后再次发起查询,重复这个过程直到满足条件或者触发超时、重试上限。整个过程完全异步,不会阻塞调用线程。
关键操作符说明
- flatMap:用于处理状态查询的结果,根据结果决定是否继续轮询
- delaySubscription:控制每次轮询的间隔时间,避免频繁请求外部系统
- repeatWhen:在满足特定条件时重复执行Mono的逻辑,适合轮询场景
- timeout:设置轮询的最大时长,避免无限轮询
- takeUntil:当满足目标状态时终止轮询流程
基础轮询示例
假设我们有一个查询外部任务状态的方法query_task_status,返回Mono<String>,状态为SUCCESS时表示任务完成,我们需要轮询直到拿到这个状态,每次轮询间隔1秒,最多轮询10次。
import reactor.core.publisher.Mono;
import reactor.util.retry.Retry;
import java.time.Duration;
public class MonoPollExample {
// 模拟查询外部系统任务状态的方法
private Mono<String> query_task_status(String taskId) {
// 实际场景中这里是调用外部接口的逻辑
return Mono.fromCallable(() -> {
// 模拟随机返回状态,实际根据外部系统返回处理
double random = Math.random();
if (random > 0.8) {
return "SUCCESS";
} else if (random > 0.6) {
return "FAILED";
} else {
return "RUNNING";
}
});
}
public Mono<String> pollTaskStatus(String taskId) {
return query_task_status(taskId)
// 判断状态是否需要继续轮询
.flatMap(status -> {
if ("SUCCESS".equals(status)) {
// 目标状态,直接返回结果
return Mono.just(status);
} else if ("FAILED".equals(status)) {
// 失败状态,抛出异常处理
return Mono.error(new RuntimeException("外部任务执行失败"));
} else {
// 非目标状态,返回空Mono触发重试
return Mono.empty();
}
})
// 当返回空Mono时,延迟1秒后重试,最多重试9次(加上初始查询共10次)
.repeatWhen(repeat -> repeat.delayElements(Duration.ofSeconds(1)).take(9))
// 设置总超时时间为15秒,超过则抛出超时异常
.timeout(Duration.ofSeconds(15))
// 当拿到SUCCESS状态时终止流
.takeUntil("SUCCESS"::equals)
// 取最后一个结果返回
.last();
}
public static void main(String[] args) {
MonoPollExample example = new MonoPollExample();
example.pollTaskStatus("task_123")
.doOnNext(status -> System.out.println("最终任务状态:" + status))
.doOnError(e -> System.out.println("轮询异常:" + e.getMessage()))
.block();
}
}
进阶场景处理
带重试策略的轮询
如果外部系统偶尔会出现查询失败的情况,可以结合retryWhen操作符增加重试逻辑,避免因为一次查询失败就终止整个轮询流程。
public Mono<String> pollTaskStatusWithRetry(String taskId) {
return query_task_status(taskId)
.flatMap(status -> {
if ("SUCCESS".equals(status)) {
return Mono.just(status);
} else if ("FAILED".equals(status)) {
return Mono.error(new RuntimeException("外部任务执行失败"));
} else {
return Mono.empty();
}
})
// 查询失败时最多重试3次,每次重试间隔500毫秒
.retryWhen(Retry.fixedDelay(3, Duration.ofMillis(500)))
.repeatWhen(repeat -> repeat.delayElements(Duration.ofSeconds(1)).take(9))
.timeout(Duration.ofSeconds(15))
.takeUntil("SUCCESS"::equals)
.last();
}
动态轮询间隔
有些场景下需要采用退避策略,比如第一次间隔1秒,第二次间隔2秒,第三次间隔4秒,这种可以通过自定义repeatWhen的逻辑实现。
public Mono<String> pollTaskStatusWithBackoff(String taskId) {
return query_task_status(taskId)
.flatMap(status -> {
if ("SUCCESS".equals(status)) {
return Mono.just(status);
} else if ("FAILED".equals(status)) {
return Mono.error(new RuntimeException("外部任务执行失败"));
} else {
return Mono.empty();
}
})
.repeatWhen(repeat -> {
// 记录重试次数,实现指数退避
return repeat.zipWith(Mono.range(1, 10), (signal, count) -> count)
.flatMap(count -> Mono.delay(Duration.ofSeconds((long) Math.pow(2, count - 1))));
})
.timeout(Duration.ofSeconds(30))
.takeUntil("SUCCESS"::equals)
.last();
}
注意事项
- 轮询间隔需要根据外部系统的承受能力设置,避免过于频繁请求导致对方服务压力过大
- 一定要设置超时时间,防止因为外部系统异常导致轮询无限进行,浪费系统资源
- 如果外部系统有查询频率限制,需要在轮询逻辑中做好限流控制
- 对于失败状态的判断要清晰,避免把可恢复的错误当成最终失败终止轮询