本节目标:说清「在事件循环线程上跑阻塞代码」的具体后果、
publishOn与subscribeOn的语义差异与误用、Schedulers.boundedElastic()的设计依据,以及用BlockHound与BlockingExecutionConfigurer检测和隔离阻塞的落地手段。
适用版本:Spring Boot 4.1.x(Java 21)
5.3 阻塞代码的隔离
WebFlux 的高并发不是靠「线程多」,而是靠「线程少、每个都不许堵」。这条纪律一旦破,整个服务的吞吐会断崖式下跌,而且现象往往和代码位置对不上:一个接口慢,全站跟着慢。本节先讲清为什么会这样,再讲三种隔离手段。
延续借阅场景:BorrowingService 需要调用一个只有同步客户端的旧风控系统——这正是最典型的「必须阻塞又必须隔离」的处境。
5.3.1 在事件循环线程上跑阻塞代码会发生什么
先看默认调度器的规模。本机在 JDK 21 上直接运行 reactor-core-3.8.7 的探针(Runtime.availableProcessors() = 10):
parallel distinct threads = 10
parallel thread name = parallel-1
boundedElastic thread name = boundedElastic-1
Schedulers.parallel() 的线程数等于 CPU 核数(本机 10 核 → parallel-1 … parallel-10)。它被 WebFlux 用作事件循环:Reactor Netty 的 I/O 事件、以及默认情况下的操作符执行,都跑在这几个线程上。
于是后果可以精确描述:
| 发生的事 | 直接后果 |
|---|---|
一个 parallel 线程被 Thread.sleep / JDBC / block() 占住 | 该线程上的所有其它请求排队等待 |
| N 个并发阻塞请求同时进来 | N 个 parallel 线程全被占满,事件循环停摆 |
| 事件循环停摆 | 连「接受新连接」「读取已就绪的数据」都做不了——包括与阻塞请求无关的接口 |
最后一行是致命的:这不是「慢接口拖慢自己」,而是一个慢接口冻结整个进程。而核数是 10,意味着只要 10 个并发阻塞请求,服务就接近假死。
5.3.2 publishOn 与 subscribeOn 的语义差异
把阻塞挪走靠这两个操作符,但它们管的是不同方向的信号:
| 操作符 | 影响的信号 | 一句话记法 |
|---|---|---|
subscribeOn(Scheduler) | 订阅信号(onSubscribe / request / cancel)向上游传播的方向 | 决定上游从哪开始跑 |
publishOn(Scheduler) | 数据信号(onNext / onComplete / onError)向下游传播的方向 | 决定下游在哪跑 |
两者的完整签名(本机核实):Flux.publishOn(Scheduler)、Flux.publishOn(Scheduler, int)、Flux.subscribeOn(Scheduler)、Flux.subscribeOn(Scheduler, boolean),以及 Mono 上的同名方法。
subscribeOn 的位置不重要,publishOn 的位置很重要。 因为 subscribeOn 作用在「订阅信号从下游往上走」的路径上——无论写在链的哪一段,它影响的是整条上游的执行线程。而 publishOn 只影响它下游的操作符。
正确用法:让阻塞调用发生在一个被 subscribeOn 或 publishOn 移出事件循环的线程上。
public Mono<RiskResult> assess(String memberId) {
return Mono.fromCallable(() -> legacyRiskClient.assess(memberId)) // 同步阻塞调用
.subscribeOn(Schedulers.boundedElastic()); // 换到弹性线程池执行
}
三种常见误用:
// 误用一:阻塞写在 fromCallable 之外 —— 调用发生在构造期,与调度无关
Mono<RiskResult> bad1 = Mono.just(legacyRiskClient.assess(memberId))
.subscribeOn(Schedulers.boundedElastic()); // 太晚了,早就阻塞完了
// 误用二:subscribeOn 写在阻塞之后,指望它把阻塞挪走
Mono<RiskResult> bad2 = Mono.fromCallable(() -> legacyRiskClient.assess(memberId))
.map(this::enrich) // 这里的执行线程由 subscribeOn 决定,没错
.subscribeOn(Schedulers.parallel()); // 但换到 parallel 等于没换,还是事件循环
// 误用三:publishOn 写在阻塞之后 —— 只影响下游,阻塞仍在原线程
Mono<RiskResult> bad3 = Mono.fromCallable(() -> legacyRiskClient.assess(memberId))
.publishOn(Schedulers.boundedElastic()); // 阻塞已经在上游线程跑完了
记住一句话:subscribeOn / publishOn 都只是「把后续执行换到哪个线程」的声明,它们不会回溯地改变已经确定的执行点。 想隔离阻塞,必须让阻塞调用本身处在被切换后的线程上。
5.3.3 Schedulers.boundedElastic():为阻塞而生的调度器
Reactor 的调度器家族(本机核实)里,只有 boundedElastic 是明确为阻塞调用设计的:
| 调度器 | 线程模型 | 用途 |
|---|---|---|
Schedulers.parallel() | 固定 = CPU 核数(本机 10) | 非阻塞的 CPU 密集计算、事件循环 |
Schedulers.boundedElastic() | 弹性增长,有上限 | 阻塞调用(JDBC、同步 HTTP、文件 I/O) |
Schedulers.single() | 单线程 | 需要串行执行的场景 |
Schedulers.immediate() | 当前线程直接执行 | 不需要切换 |
Schedulers.fromExecutor(Executor) | 包装已有线程池 | 复用业务线程池 |
boundedElastic 的两个关键默认值(本机核实常量 + 探针实测):
| 常量 | 值 | 含义 |
|---|---|---|
Schedulers.DEFAULT_BOUNDED_ELASTIC_SIZE | 100 | 线程数上限(不是每核 100,是总量 100) |
Schedulers.DEFAULT_BOUNDED_ELASTIC_QUEUESIZE | 100000 | 每个线程的排队任务上限 |
Schedulers.DEFAULT_BOUNDED_ELASTIC_ON_VIRTUAL_THREADS | false | 默认不基于虚拟线程 |
探针实测印证了线程上限:同时提交 200 个各 sleep(200ms) 的阻塞任务,观察到的不同线程数为 100(即达到上限后开始排队,而不是继续新建线程)。
为什么这个设计合理:
- 有上限:阻塞任务再多也不会无限建线程把机器压垮,超出部分排队。
- 可弹性增长:空闲时线程会被回收,低峰期不占资源。
- 与
parallel隔离:事件循环线程永远不被阻塞任务占用,这是整套设计的核心。
它「为阻塞而生」的另一面是:不要在 boundedElastic 上跑 CPU 密集任务。CPU 密集任务会长时间占住弹性线程,既拖慢其它阻塞任务,又浪费了「弹性」这个特性——这类任务应交给 parallel。
5.3.4 BlockHound:把阻塞调用抓出来
光靠人眼审查很难找全阻塞点,BlockHound 用字节码插桩在运行时检测「非阻塞线程上的阻塞调用」。它不在 reactor-core 里,而在独立的 io.projectreactor.tools:blockhound(本机已核实类结构,版本 1.0.13.RELEASE):
public class reactor.blockhound.BlockHound {
public static reactor.blockhound.BlockHound$Builder builder();
public static void install(reactor.blockhound.integration.BlockHoundIntegration... integrations);
public static void premain(java.lang.String, java.lang.instrument.Instrumentation);
}
BlockHound$Builder 的常用方法(本机核实):
| 方法 | 作用 |
|---|---|
allowBlockingCallsInside(String className, String methodName) | 放行白名单内的阻塞调用 |
disallowBlockingCallsInside(String, String) | 反向:禁止 |
markAsBlocking(Class, String, String) | 把某个方法标记为阻塞 |
blockingMethodCallback(Consumer<BlockingMethod>) | 自定义命中回调 |
nonBlockingThreadPredicate(Function<...>) | 自定义「哪些线程算非阻塞」 |
addDynamicThreadPredicate(Predicate<Thread>) | 动态标记线程 |
with(BlockHoundIntegration) / loadIntegrations(...) | 装载集成 |
install() | 生效 |
在测试里启用(需在 JVM 参数加 -XX:+AllowRedefinitionToAddDeleteMethods,这是 BlockHound 的硬性要求):
@BeforeAll
static void installBlockHound() {
BlockHound.builder()
.allowBlockingCallsInside("java.io.PrintStream", "write") // 放行日志
.install();
}
命中时的报错形如(示例输出——本机未在完整 WebFlux 应用内运行 BlockHound,以下为按官方文档格式书写的示例):
reactor.blockhound.BlockingOperationError: Blocking call! java.io.FileInputStream#readBytes
at java.base/java.io.FileInputStream.readBytes(FileInputStream.java)
at com.example.borrowing.LegacyRiskClient.assess(LegacyRiskClient.java:42)
at com.example.borrowing.BorrowingService.assess(BorrowingService.java:31)
...
BlockHound 只适合测试期与预发环境——它靠插桩,有性能开销,不应长期开在生产。它的价值是把「哪一行在事件循环上阻塞」从猜测变成确定的栈。
5.3.5 WebFlux 的阻塞执行隔离:BlockingExecutionConfigurer
除了在代码里手动 subscribeOn,Spring 还提供了一层声明式隔离。WebFluxConfigurer(本机核实)里有一个默认方法:
default void configureBlockingExecution(org.springframework.web.reactive.config.BlockingExecutionConfigurer configurer);
BlockingExecutionConfigurer(本机核实)只有两个配置点:
| 方法 | 作用 |
|---|---|
setExecutor(AsyncTaskExecutor) | 指定跑阻塞控制器方法的执行器 |
setControllerMethodPredicate(Predicate<HandlerMethod>) | 指定「哪些控制器方法算阻塞」 |
配置示例:
@Configuration
public class WebFluxBlockingConfig implements WebFluxConfigurer {
@Override
public void configureBlockingExecution(BlockingExecutionConfigurer configurer) {
configurer.setControllerMethodPredicate(handlerMethod ->
handlerMethod.hasMethodAnnotation(Blocking.class));
// 默认执行器即 boundedElastic;如需自定义可 setExecutor(...)
}
}
配合一个自定义注解:
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface Blocking {
}
给必须阻塞的控制器方法打上 @Blocking,框架就会把这些方法的执行调度到阻塞执行器上,而不是事件循环——把「哪些方法会阻塞」从口头约定变成可配置、可审计的声明。
顺带一提,Boot 自己的 Actuator 也用了同样的思路:org.springframework.boot.webflux.actuate.endpoint.web.AbstractWebFluxEndpointHandlerMapping 内部有一个 ElasticSchedulerInvoker(本机核实),专门把端点调用切到弹性调度器上,避免管理端点的阻塞逻辑影响主事件循环。
5.3.6 排查清单
| 现象 | 优先怀疑 | 处理 |
|---|---|---|
全站接口一起变慢,线程名都是 parallel-* | 事件循环被阻塞 | 用 BlockHound 定位,把阻塞切到 boundedElastic |
加了 subscribeOn 仍然卡 | 阻塞写在 fromCallable 之外 | 把阻塞移进 Mono.fromCallable / Flux.defer |
换了 publishOn 没效果 | 只影响下游,阻塞在上游 | 改用 subscribeOn 或把阻塞移到切换点之后 |
boundedElastic 任务大量排队 | 阻塞任务量超过 100 线程上限 | 提高上限(newBoundedElastic)或减少阻塞、加缓存 |
| 阻塞任务把 CPU 打满 | 在弹性池上跑了 CPU 密集任务 | 换到 Schedulers.parallel() |
5.3.7 本机可以做的验证
export JAVA_HOME=/tmp/springboot_book/jdk-21.0.12.1+1/Contents/Home
J=/tmp/springboot_book/jars
# 1) 调度器常量与工厂方法
"$JAVA_HOME/bin/javap" -cp "$J/reactor-core-3.8.7.jar" reactor.core.scheduler.Schedulers
# 2) 打印默认常量值(本机实测:SIZE=100、QUEUESIZE=100000、ON_VIRTUAL_THREADS=false)
cat > /tmp/Probe.java <<'EOF'
import reactor.core.scheduler.Schedulers;
public class Probe {
public static void main(String[] a) {
System.out.println("SIZE=" + Schedulers.DEFAULT_BOUNDED_ELASTIC_SIZE);
System.out.println("QUEUESIZE=" + Schedulers.DEFAULT_BOUNDED_ELASTIC_QUEUESIZE);
System.out.println("ON_VT=" + Schedulers.DEFAULT_BOUNDED_ELASTIC_ON_VIRTUAL_THREADS);
}
}
EOF
"$JAVA_HOME/bin/javac" -cp "$J/reactor-core-3.8.7.jar" -d /tmp /tmp/Probe.java
"$JAVA_HOME/bin/java" -cp "$J/reactor-core-3.8.7.jar:/tmp" Probe
# 3) BlockHound 的类与配置项
"$JAVA_HOME/bin/javap" -cp "$J/blockhound-1.0.13.RELEASE.jar" reactor.blockhound.BlockHound
"$JAVA_HOME/bin/javap" -cp "$J/blockhound-1.0.13.RELEASE.jar" 'reactor.blockhound.BlockHound$Builder'
# 4) WebFlux 的阻塞执行配置点
"$JAVA_HOME/bin/javap" -cp "$J/spring-webflux-7.0.9.jar" \
org.springframework.web.reactive.config.BlockingExecutionConfigurer
线程上限那条结论,把 5.3.3 描述的探针(提交 200 个 sleep 任务、统计不同线程数)跑一遍即可复现——这是本机实测,不是示例输出。而 BlockHound 在完整 WebFlux 应用里的命中栈属示例输出,本机没有运行该场景。
5.3.8 知道之后能做什么
给旧同步客户端加壳。 面对只有同步 API 的风控、支付、老 ERP,统一用 Mono.fromCallable(...).subscribeOn(Schedulers.boundedElastic()) 包一层,把它变成不污染事件循环的 Mono。
把阻塞点写进代码规约。 用 @Blocking 注解 + BlockingExecutionConfigurer 把「哪些方法允许阻塞」变成声明,配合 CI 里的 BlockHound 测试,防止新代码悄悄把阻塞带回事件循环。
调容量时先分清瓶颈类型。 吞吐上不去时,先判断是「阻塞任务把弹性池占满」(调 boundedElastic 上限或加缓存)还是「CPU 密集任务挤占 parallel」(改调度器)——两者的处方完全相反。
小结
- WebFlux 的事件循环线程数等于 CPU 核数(本机 10 核 → 10 条
parallel-*);一条被阻塞,该线程上的所有请求一起等,核数用尽则全站停摆。 subscribeOn影响上游执行起点、publishOn影响下游;两者都不会回溯地改变已确定的执行点,把阻塞写在Mono.just(...)里再用subscribeOn是无效的。Schedulers.boundedElastic()专为阻塞设计:默认线程上限 100(本机实测)、每线程排队上限 100000,与事件循环隔离;不要在它上面跑 CPU 密集任务。BlockHound通过字节码插桩在运行时抓出「非阻塞线程上的阻塞调用」,适合测试与预发,不宜长期开在生产。WebFluxConfigurer.configureBlockingExecution+BlockingExecutionConfigurer能把「哪些控制器方法阻塞」声明化,由框架调度到阻塞执行器。
本章到此结束:5.1 讲清发布者与背压,5.2 讲清请求链,5.3 讲清唯一必须守住的纪律——不要阻塞事件循环。下一章转向数据访问层:EntityManager 与持久化上下文在框架内部是怎么被管理的。
阅读导航:上一节:5.2 WebFlux 请求处理 · 下一节:6.1 EntityManager 与持久化上下文 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。