《Spring Boot 高级》5.3 阻塞代码的隔离

解释在事件循环线程上执行阻塞调用为何会拖垮整个服务,厘清 publishOn 与 subscribeOn 的语义差异与常见误用,说明 Schedulers.boundedElastic() 为何专为阻塞设计(本机实测默认 100 线程上限),并给出 BlockHound 检测与 WebFlux 阻塞执行隔离的配置方法。

本节目标:说清「在事件循环线程上跑阻塞代码」的具体后果、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_SIZE100线程数上限(不是每核 100,是总量 100)
Schedulers.DEFAULT_BOUNDED_ELASTIC_QUEUESIZE100000每个线程的排队任务上限
Schedulers.DEFAULT_BOUNDED_ELASTIC_ON_VIRTUAL_THREADSfalse默认不基于虚拟线程

探针实测印证了线程上限:同时提交 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 与持久化上下文 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java」更多文章

  1. 《Spring Boot 入门》18.3 打包与运行
  2. 《Spring Boot 入门》18.2 实现
  3. 《Spring Boot 入门》18.1 需求与设计