Reactor中timeout算子的计时范围及失效场景解析
Redis反应式代码Timeout问题及解答
问题背景
近期一段Redis反应式代码出现timeout问题,代码示例如下:
redisOps.entries(key).map(...).map(...).switchIfEmpty().timeout();
为确认超时是否发生在Redis查询阶段,计划将timeout()置于entries()之后,避免map操作的耗时被计入超时监控。编写测试代码后,提出以下疑问:
疑问1
timeout()是否仅监控其之前的数据发射耗时,后续操作不计入?例如以下代码中,仅a和b操作的耗时被timeout()监控,c和d不计入:
Mono mono = foo(); mono.a().b().timeout().c().d();
疑问2
为何以下代码中的timeout()不生效?
Mono.just("good luck") .map(s -> { try { TimeUnit.SECONDS.sleep(1); } catch (InterruptedException e) { e.printStackTrace(); } return "timeout 1"; }) .map(s->{ try { TimeUnit.SECONDS.sleep(1); } catch (InterruptedException e) { e.printStackTrace(); } return "timeout 2"; }) .timeout(Duration.ofSeconds(1)) .subscribe(System.out::println);
疑问3
与问题2相比,为何以下代码中的timeout()生效?
public static void main(String[] args) { Mono.fromCallable(() -> { log.info("begin 1"); try { TimeUnit.SECONDS.sleep(2); } catch (InterruptedException e) { e.printStackTrace(); } log.info("return 1"); return "good luck 1"; }) .timeout(Duration.ofMillis(1500l)) .map((s) -> { log.info("begin 2"); try { TimeUnit.SECONDS.sleep(1); } catch (InterruptedException e) { e.printStackTrace(); } log.info("return 2"); return "good luck 2"; }) .subscribe(System.out::println); }
解答
疑问1解答
是的,timeout()仅监控其上游操作链的数据发射耗时——从订阅开始,到timeout()之前的整个操作链完成并发射数据的时间。后续的c()和d()属于timeout()的下游操作,它们的执行耗时不会被计入timeout()的监控范围。
本质上,timeout()是针对它前面的Mono/Flux序列:如果上游在指定时间内没有发射第一个数据(Mono场景)或下一个数据(Flux场景),就会触发超时。下游操作的执行和timeout()的逻辑完全无关。
疑问2解答
问题出在Mono.just()的特性和map的执行时机:
Mono.just("good luck")是立即完成的Mono,订阅后会立刻发射数据。map操作是同步执行的,在订阅阶段,两个map里的sleep会被连续执行(总共2秒),但这整个过程是在timeout()开始监控之前就完成了。- 当
timeout()被调用时,上游已经完成了数据发射,timeout()没有机会触发——它只会监控订阅后上游是否在指定时间内发射数据,而这里上游在timeout()生效前就已经完成了。
疑问3解答
Mono.fromCallable()和Mono.just()的执行逻辑有本质区别:
fromCallable()中的逻辑是订阅后延迟执行的(在订阅触发后,由Reactor调度器执行)。- 订阅这个Mono时,
timeout()会立刻开始计时,同时fromCallable()里的sleep(2秒)开始执行。因为2秒超过了timeout()设置的1500毫秒,所以在fromCallable()返回数据之前,timeout()就触发了超时。 - 这里的
timeout()在订阅后立刻生效,正好覆盖了fromCallable()的执行耗时,因此能正确触发超时。
内容的提问来源于stack exchange,提问作者Tiina
相关产品推荐
相关产品推荐

