You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring WebFlux中Flux在WebClient与Flux.just场景下行为差异的原因咨询

问题拆解:为什么Flux.just的subscribe会阻塞主线程?

这是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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.29 20:42:28