Mono.just()在元素生成阻塞时为何阻塞?Reactor使用疑问
你遇到的问题核心是没搞清楚Mono.just()的求值时机,以及Reactor非阻塞特性的前提条件,下面逐一拆解:
测试1符合预期的原因
第一个测试里,Mono.just("test")直接发射字符串,map里的阻塞sleep能表现出非阻塞,本质是你应该用了subscribeOn或publishOn把map中的阻塞逻辑调度到了单独的线程池,让主线程不会被阻塞,这是正确利用Reactor异步调度的方式。
测试2阻塞的根本错误
第二个测试用Mono.just(getString())时,Mono.just()的参数是在调用它的当前线程(通常是主线程)立即执行的。也就是说,你的getString()方法里的sleep,在你创建这个Mono对象的瞬间就同步执行了,完全没等到Reactor的异步调度环节。
代码等价于:
// 这里直接阻塞当前线程,和Reactor无关 String syncResult = getString(); // 只是把已经拿到的结果包装成Mono Mono<String> mono = Mono.just(syncResult);
这种写法相当于把阻塞逻辑放在了Reactor的异步流程之外,自然会导致整个调用阻塞。
测试3非阻塞的逻辑
第三个测试用Flux.create实现非阻塞,是因为create的回调逻辑是在订阅后才执行的,你可以把sleep这类阻塞操作放到回调里,再配合调度器(比如Schedulers.boundedElastic())把逻辑放到异步线程,避免阻塞主线程,这才是符合Reactor非阻塞设计的用法。
修正测试2的正确写法
要让第二个测试实现非阻塞,你需要把阻塞的getString()包装成延迟求值的操作,用Mono.fromSupplier()或Mono.defer()即可:
方案1:使用Mono.fromSupplier()
Mono.fromSupplier(this::getString) .subscribeOn(Schedulers.boundedElastic()) // 将阻塞操作交给专门的线程池处理 .subscribe(res -> System.out.println(res));
fromSupplier会把getString()的执行推迟到订阅阶段,再通过subscribeOn指定线程池,确保阻塞逻辑不会卡主线程。
方案2:使用Mono.defer()
Mono.defer(() -> Mono.just(getString())) .subscribeOn(Schedulers.boundedElastic()) .subscribe(res -> System.out.println(res));
defer会延迟整个Mono的创建过程,直到订阅发生,同样能把阻塞逻辑放到异步线程中执行。
关键概念梳理
Reactor的非阻塞不是“套个Mono/Flux就自动生效”,核心是两点:
- 所有阻塞操作必须放到专门的调度线程池(比如
boundedElastic),不能在订阅线程执行 - 区分立即求值(如
Mono.just())和延迟求值(如fromSupplier()、defer())的操作,阻塞逻辑必须用延迟求值的方式包装
内容的提问来源于stack exchange,提问作者Rionash

