Spring WebFlux中Flux在WebClient与Flux.just场景下行为差异的原因咨询
这是Reactor入门阶段非常典型的疑问,核心在于你忽略了Reactor操作的线程模型差异,我来给你一步步理清楚:
先看两个例子的本质区别
第一个WebClient场景
WebClient是Reactor封装的异步HTTP客户端,它发起请求时,内部会自动把请求和后续的回调逻辑切换到Reactor的异步线程池(比如reactor-http-nio-*系列线程)中执行。所以当你调用tweetFlux.subscribe(...)时,订阅逻辑并没有在Controller的主线程里执行,主线程会直接跳过subscribe继续往下走,自然会先打印Exiting NON-BLOCKING Controller!,等异步线程处理完请求后才会打印tweet内容。
第二个Flux.just场景
Flux.just(1,2,3,4,5)创建的是一个同步的、基于当前线程执行的冷序列。Reactor的默认行为是:如果没有显式指定线程调度器,整个订阅和消费过程都会在调用subscribe()的线程(也就是Controller处理请求的主线程)上同步执行。
也就是说,当你调用f.subscribe(...)时,主线程会立刻开始遍历Flux的所有元素,逐个执行回调里的逻辑(包括Thread.sleep(2000L)),等5个元素都处理完(总共10秒),才会继续执行后面的log.info("doing something else")——这就是为什么你看到的输出顺序和预期相反。
怎么让第二个例子也实现非阻塞?
你只需要显式指定订阅时使用异步线程池,用subscribeOn()操作符切换线程即可:
@GetMapping("/test") public void doSomething(){ System.out.println("i am here"); Flux<Integer> f= Flux.just(1,2,3,4,5); f.subscribeOn(Schedulers.boundedElastic()) // 指定用弹性线程池执行订阅和消费 .subscribe(consumer->{ try { log.info("consuming"); Thread.sleep(2000L); } catch (InterruptedException e) { e.printStackTrace(); } log.info(String.valueOf(consumer)); }); log.info("doing something else"); }
修改后的输出就会符合你的预期:
i am here doing something else consuming 1 consuming 2 consuming 3 consuming 4 consuming 5
关键总结
Reactor的“非阻塞”不是凭空实现的:
- 如果数据源本身是异步的(比如WebClient的HTTP请求、数据库的异步驱动),Reactor会自动在异步线程处理;
- 如果是同步数据源(比如Flux.just、Flux.fromIterable处理内存集合),必须手动通过
subscribeOn或publishOn指定异步线程,才能避免阻塞当前线程。
内容的提问来源于stack exchange,提问作者gooner

