上一章结束时,TaskHub 已经有了两张成熟的面孔:程序可以通过 /api/tasks 返回 JSON,也可以通过 /tasks 渲染任务看板。两条路径最终都复用同一个 TaskService,再由 Spring Data JPA 访问数据库。这套 MVC + JPA 结构没有过时,也没有因为这一章出现就需要推倒重来。
现在产品多了一个不太一样的需求:运营同事打开看板后,希望任务状态一变化,页面就能立刻收到消息,而不是每隔几秒发一次查询。TaskHub 还要向另一个成员目录服务查询负责人姓名。前者是一条可能持续很久的连接,后者的大部分时间花在等待网络。如果仍然让一个请求线程从头等到尾,当然也能做,但当同时在线的连接很多时,线程会有相当长的时间只是在等。
这正是我们讨论 Project Reactor 和 Spring WebFlux 的理由。不过先把结论说在前面:响应式不等于更快,WebFlux 也不是 Spring MVC 的升级版。 它换了一套并发模型,用少量事件循环线程照看大量处于等待状态的 I/O。只有当整条调用链大体上都能非阻塞时,这种模型才真正合算。
因此,我们不会把现有 TaskHub 原地改造成 WebFlux,更不会把阻塞式 JPA 塞进事件循环。主应用继续在 8080 端口负责 CRUD 和管理页面;这一章建立一个独立的 taskhub-reactive 应用,运行在 8081 端口,专门负责持续任务动态,并借它讲清 Reactor 到底在做什么。
第一次接触 WebFlux 时,人很容易被“高并发”“非阻塞”几个词推着走,好像只要依赖换成 WebFlux,吞吐量就会自动提高。我们先把 TaskHub 的需求拆开看。
普通的任务增删改查要访问 JPA。JPA 和底层 JDBC 的契约是阻塞式的:数据库结果没有回来,执行查询的线程就要等着。现有 MVC 应用的线程池本来就是按“业务代码可能阻塞”这个前提设计的,它很适合这条路径。对一个规模正常、以数据库 CRUD 为主的系统,清楚的命令式代码往往更容易开发和排错。
持续事件则不同。浏览器订阅 /api/task-events 后,连接可能保持几分钟甚至几小时,绝大多数时间没有数据可写。成员目录查询也类似,CPU 真正处理 JSON 的时间很短,主要延迟来自网络。这样的 I/O 等待型场景,才有可能从非阻塞模型中获益。
这里值得再拆掉三个常见误解。
第一个误解是“异步就是单次请求更快”。一次成员查询原来要 300 毫秒,改用 WebClient 后它大概率还是要 300 毫秒。变化在于等待这 300 毫秒的方式:阻塞模型让一个线程停住,非阻塞模型登记后续动作并释放事件循环线程。只有要同时等待多个彼此独立的上游时,合理并发才可能缩短这一组调用的总等待时间。
第二个误解是“线程少就没有容量上限”。网络连接、连接池、内存缓冲、下游服务和 CPU 都有上限。WebFlux 让等待资源的方式更省,不会把有限资源变成无限。没有并发控制的 flatMap、没有边界的缓冲区和无限重试,同样能把响应式服务拖垮,只是故障形态和线程池耗尽不太一样。
第三个误解是“一个项目只能选一种”。MVC 应用可以单独使用 WebClient 调用远程服务,也能返回某些响应式类型;一个系统里的不同服务也可以按依赖特点分别选择 MVC 或 WebFlux。真正需要避免的,是在同一条请求路径里来回切换思维模型,却没有说明阻塞发生在哪里、事务归谁管理、错误如何回到客户端。
选择 WebFlux 时,先列出请求路径上的每个依赖:HTTP 客户端、数据库驱动、文件访问、第三方 SDK。只要核心步骤仍然大量阻塞,就很难得到端到端非阻塞的收益。框架名称不是判断依据,依赖的真实行为才是。
我们这一章选择“两个应用并存”,并不是为了炫技,而是让边界一眼可见:taskhub 继续使用 MVC + JPA,taskhub-reactive 只处理适合流式和非阻塞的工作。以后如果响应式部分收益不明显,删掉或合并它也不会动摇主业务。
做技术选型时,可以把问题落到几项能测量的事实上:预计同时保持多少连接、每次请求要等待几个上游、延迟分布是否经常出现长尾、阻塞依赖占多少、团队能否定位异步调用栈、现有 MVC 在目标负载下究竟哪里先到瓶颈。没有这些事实,仅凭“访问量以后可能很大”就重写框架,通常会把确定的复杂度换成还没被证明的收益。
响应式应用的坐标是 com.welearn:taskhub-reactive:0.0.1-SNAPSHOT,主包名是 com.welearn.taskhub.reactive。课程统一使用 Spring Boot 4.1.0 和 Java 25。它的 pom.xml 中保留 WebFlux 和测试所需的依赖:
<properties>
<java.version>25</java.version>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<
这里有两个容易踩的坑。
第一个坑,是把 spring-boot-starter-webmvc 和 spring-boot-starter-webflux 一股脑放进主项目,然后以为控制器返回 Mono 就会自动变成 WebFlux。两套框架确实可以出现在同一个依赖图里,但 Spring Boot 同时看到它们时通常会选择 MVC 应用类型。更关键的是,即使 MVC 控制器能返回响应式类型,底层的请求模型也不会因此变成 WebFlux。判断当前跑的是哪套 Web 栈,不能只看方法返回值。
第二个坑,是看到 Flux 就顺手把 JPA Repository 注入进来。返回类型可以包装,阻塞行为却不会被包装掉。Mono.fromCallable(() -> jpaRepository.findAll()) 里面的 findAll() 仍然会让某个线程等待。后面我们会讲临时隔离办法,但这个独立应用先保持边界干净。
启动类仍然很普通:
package com.welearn.taskhub.reactive;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class TaskHubReactiveApplication {
public static void main(String[] args) {
SpringApplication.run(TaskHubReactiveApplication.class, args);
}
}在 application.yml 中固定应用名、端口和需要暴露的管理端点:
spring:
application:
name: taskhub-reactive
server:
port: 8081
management:
endpoints:
web:
exposure:
include: health,info用于跑通 SSE 的项目先保持克制,只放本章真正执行的角色:
com.welearn.taskhub.reactive
├── TaskHubReactiveApplication
├── TaskEvent
├── TaskEventService
├── TaskEventController
└── ResilientTaskLookup文件数量还很少,先放在根包更便于跟读。后面把成员客户端、R2DBC 映射和更多任务逻辑真正加入应用时,再按 event、member、task 拆包。分包不是为了让文件数量看起来“像企业项目”,而是让变更理由分开:SSE 帧格式变化不应改成员客户端,成员服务换地址也不该改任务事件源。还没有这个复杂度时,过早造一排空目录只会增加噪声。
@SpringBootApplication 在这里仍然承担配置类、自动配置入口和组件扫描起点三个角色。不同的是,类路径上有 WebFlux starter 后,自动配置会建立响应式 Web 基础设施。控制器返回的 Mono 或 Flux 会交给 WebFlux 处理,默认的 Reactor Netty 服务器用非阻塞网络 I/O 读写请求和响应。
运行 ./mvnw spring-boot:run 后,启动日志明确显示当前不是 Servlet Web 服务器:
Netty started on port 8081 (http)
Started TaskHubReactiveApplication in 1.004 seconds第一行确认 Reactor Netty 正在监听 8081;第二行说明 Spring 容器已经完成 Bean 创建和 WebFlux 基础设施启动。日志里看见 Netty 只证明服务器栈选对了,仍不能证明你的每个业务依赖都非阻塞,后者要逐段检查。
我们先从两个最常见的类型开始:
Mono<T> 表示一条最终可能发出 0 个或 1 个 T 的序列。Flux<T> 表示一条最终可能发出 0 到多个 T 的序列,序列可以结束,也可以像事件流一样长期存在。“0 个”不是文字游戏。按 id 查询任务时,找到就发出一个任务,没找到可以正常结束但不发出值,所以返回 Mono<TaskSnapshot> 很合适。删除操作如果只关心完成与失败,可以用 Mono<Void>。任务事件可能不断到达,则用 Flux<TaskEvent>。

为了不让数据库问题干扰 Reactor 的第一步,我们先准备一个内存目录:
package com.welearn.taskhub.reactive.task;
public record TaskSnapshot(
long id,
String title,
String status,
int priority,
long assigneeId) {
}package com.welearn.taskhub.reactive.task;
import org.springframework.stereotype.Component;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.util.List;
@Component
public class TaskCatalog {
private final List<TaskSnapshot> tasks = List.of(
new TaskSnapshot(1, "补上接口测试", "TODO", 5, 101),
findAll() 没有把 List 改成一个更时髦的容器。它返回的是一个 Publisher,也就是一份“有人订阅时,应该怎样产生和传递任务”的说明。findById() 先筛选,再通过 next() 取第一个元素,所以类型从 Flux<TaskSnapshot> 变成 Mono<TaskSnapshot>。找不到时,它会得到一个空的 Mono,而不是 null。
这和 Java Stream 有一点相似:两者都能 map、filter,也都把一串变换声明出来。但 Java Stream 主要处理同步拉取的集合数据;Reactor 还要面对值晚点才到、序列可能没有终点、错误和取消也是信号、下游需要控制请求量等情况。不要因为操作符名字相同,就把 Flux 当成“异步版 List”。
null 也不属于响应式数据流的正常元素。一个 Mono 没有结果时,应返回 Mono.empty();一个可能为 null 的现成值可以用 Mono.justOrEmpty(value);异步失败则用 Mono.error(error)。把“没有数据”和“处理失败”分开后,控制器才能准确决定是 404、204,还是 5xx。
Reactor 序列传递的不是只有值。一次正常订阅大致会经历这些信号:
onSubscribe 交出订阅关系。request(n) 表示自己还能接收多少个元素。onNext 传递元素。onComplete 正常结束,或以一次 onError 失败结束。onComplete 和 onError 是互斥的终止信号。流一旦失败,原来的上游不会越过错误继续发后面的元素。背压也从这里进入:订阅者不是被动地无限接收,而是用 request(n) 表达需求。我们这一章先建立这个心智模型,第十二章再用可控请求量、缓冲和取消把它拆得更细。
信号的方向也很重要。订阅建立和 request(n) 从下游往上游走,数据、完成与错误从上游往下游走。假设管道是“数据源 → filter → map → HTTP 编码器”,HTTP 写入端要先表达自己能接收数据,上游才会发;元素经过 filter 和 map 后到达编码器。操作符不是按代码行从上到下主动跑一遍,而是每个环节都参与同一条订阅链。
Mono 并不表示“代码一定在后台线程执行”,Flux 也不表示“元素会并行处理”。它们描述的是序列的基数和信号协议。线程在哪里执行,由信号来源、服务器运行时和调度器共同决定。
看一段最小的操作链:
Flux<String> importantTitles = taskCatalog.findAll()
.filter(task -> task.priority() >= 4)
.map(TaskSnapshot::title)
.doOnSubscribe(subscription -> System.out.println("开始订阅任务"))
.doOnNext(title -> System.out.println("收到标题:" + title))
.
如果代码到这里就结束,输出只有:
操作链已经装配filter 没有筛选任何东西,map 也没有拿到标题。这一阶段叫装配:我们只是把数据源和操作符连成一条管道。真正有订阅者之后,订阅信号才会从下游向上游建立关系,元素再从上游向下游流动。
为了在普通 Java 示例里看见它,可以暂时手动订阅:
importantTitles.subscribe();随后输出是:
开始订阅任务
收到标题:补上接口测试
收到标题:整理部署清单
任务流完成
但到了 WebFlux 控制器和业务服务里,通常不要自己调用 subscribe()。控制器只需要把 Publisher 返回给框架:
@GetMapping("/api/reactive-tasks")
public Flux<TaskSnapshot> list() {
return taskCatalog.findAll();
}WebFlux 在准备写 HTTP 响应时会成为订阅者。它请求数据、把元素交给编码器、写入网络缓冲区,并在客户端断开时发出取消。业务代码自己 subscribe(),相当于另开了一条和 HTTP 生命周期脱节的支线:请求可能早已返回,错误也可能没人处理,测试还很难判断它何时结束。
如果想亲眼看见“订阅者决定要多少”,可以在一个纯演示中使用 BaseSubscriber。下面的订阅者一次只要一个任务,处理完再申请下一个:
taskCatalog.findAll()
.doOnRequest(n -> System.out.println("上游收到需求:" + n))
.subscribe(new BaseSubscriber<>() {
@Override
protected void hookOnSubscribe(Subscription subscription) {
request(1);
}
@Override
protected void hookOnNext(TaskSnapshot task) {
输出的交替顺序很直观:
上游收到需求:1
处理:补上接口测试
上游收到需求:1
处理:整理部署清单
上游收到需求:1
处理:核对健康检查
上游收到需求:1
全部完成这不表示实际 HTTP 响应总会逐个调用 request(1)。WebFlux 和操作符会根据预取、网络写入能力等因素批量申请。这个例子只用来说明协议:下游有表达承载能力的通道。若数据源完全无法放慢,就必须明确选择缓冲、丢弃、只保留最新值或失败,而不是假设背压会凭空控制所有外部世界。
map 做同步的一对一转换。把 TaskSnapshot 变成标题字符串,输入一个元素,输出一个元素,过程不需要等待另一个 Publisher:
Flux<String> titles = taskCatalog.findAll()
.map(TaskSnapshot::title);filter 决定一个元素是否继续向下游走:
Flux<TaskSnapshot> unfinished = taskCatalog.findAll()
.filter(task -> !task.status().equals("DONE"));flatMap 用在“一个元素要启动另一个异步步骤”的地方。假设每个任务都要去成员目录查负责人,findMember 返回的是 Mono<MemberProfile>,那么直接 map 会得到 Flux<Mono<MemberProfile>>。flatMap 会订阅这些内部 Publisher,再把结果展平成一条流:
Flux<TaskCard> cards = taskCatalog.findAll()
.flatMap(task -> memberClient.findMember(task.assigneeId())
.map(member -> TaskCard.from(task, member)));因为多个内部调用可能同时进行,flatMap 的完成顺序不一定等于任务输入顺序。必须保持顺序时,可以用 concatMap,它会等前一个内部 Publisher 完成后再处理下一个;想保序又想保留一定并发,可以使用 flatMapSequential。这里不是谁“更高级”,而是吞吐量、延迟和顺序语义之间的取舍。
还要留意 doOnNext、doOnError、doFinally 这一类以 doOn 开头的操作符。它们用于日志、指标或轻量观察,不改变流里的值。把核心写操作藏在 doOnNext 中会让业务含义变得模糊,而且副作用可能因重试或二次订阅重复执行。需要“保存完成后再继续”的动作,应让保存本身返回 Publisher,再用 flatMap、then 等操作符明确组合,而不是把关键步骤伪装成调试钩子。
“订阅后才执行”适用于常见的冷 Publisher,但响应式世界还有另一类源。理解冷流和热流,能解释不少看起来像重复请求或漏事件的问题。
TaskCatalog.findAll() 是冷流。每个订阅者都会从第一个任务开始,独立执行筛选和转换。WebClient 发出的 HTTP 请求通常也是冷的:同一个 Mono 被订阅两次,默认会发出两次网络请求。
下面的代码把“读取发生的时机”写得更明显:
Flux<TaskSnapshot> snapshots = Flux.defer(() -> {
System.out.println("重新读取任务快照");
return taskCatalog.findAll();
});
snapshots.take(1).subscribe(task -> System.out.println("甲:" + task.title()));
snapshots.take(1).subscribe(task ->输出中“重新读取任务快照”会出现两次:
重新读取任务快照
甲:补上接口测试
重新读取任务快照
乙:补上接口测试defer 适合把有副作用或依赖当前状态的创建逻辑推迟到订阅时。对比下面两行就更清楚:
Mono<TaskSnapshot> eager = Mono.just(loadTaskNow());
Mono<TaskSnapshot> lazy = Mono.defer(() -> Mono.just(loadTaskNow()));构造 eager 时,loadTaskNow() 已经执行;构造 lazy 时没有,订阅才执行。这里真正造成提前执行的是 Java 先计算了 Mono.just(...) 的参数,不是 just 在背后偷偷订阅。
这一区别会直接影响时间、认证信息和数据库读取。若组装时捕获了 Instant.now(),所有订阅者可能看到同一个旧时间;若用 defer 在订阅时创建,才会得到各自订阅发生时的值。再比如 Web 请求里的用户身份通常存在响应式 Context 中,它也是沿订阅关系传递,不应在应用启动阶段就读取。
任务动态更像广播。任务在 10:00:00 变更,10:00:05 才连上的浏览器通常只接收之后的新事件,不会让整个世界为它重新演一遍。这就是热流的典型语义:源可以独立于某个订阅者存在,多个订阅者共享正在发生的数据。

“热”不等于“总会保存历史”。多播源可能只发给当前在线订阅者;需要新订阅者先看最近一条,可以使用有界 replay;需要可靠补发全部业务事件,就要有消息代理、事件表或其他持久化机制。不能看到 cache() 或 replay 好用,就给无限事件流加一个没有上限的缓存,那会把内存当成永远增长的事件仓库。
把冷流共享成热流也要先问业务问题。两个订阅者是否允许共用同一次远程请求?结果能缓存多久?第一次请求失败时,错误要不要被缓存?最后一个订阅者离开后,上游是否应该取消?这些答案决定使用 share、有界 replay、带时效的 cache,还是保留每次独立执行。操作符只能实现策略,不能替业务决定策略。
Spring MVC 的默认假设是“应用代码可能阻塞”。Servlet 容器准备一个较大的线程池:请求进来拿一个线程,查数据库或调用远程服务时,这个线程可以停在那里等待。并发量上升后,需要更多线程吸收等待,线程栈和上下文切换也会增加资源消耗。
WebFlux 的默认假设恰好相反:应用代码不阻塞当前线程。Reactor Netty 使用少量、固定数量的事件循环线程处理网络事件。一个远程调用发出后,线程注册好回调就可以去处理别的连接;响应到达时,运行时再继续推动对应的操作链。

这解释了为什么 WebFlux 能用较少线程照看大量等待中的连接,也解释了它最怕什么。假设事件循环里执行了下面的代码:
@GetMapping("/bad")
public Mono<TaskSnapshot> bad() {
return Mono.fromCallable(() -> {
Thread.sleep(2_000);
return loadFromJdbc();
});
}Mono.fromCallable 只把调用变成延迟执行,并没有让 Thread.sleep 或 JDBC 变成非阻塞。订阅发生后,如果它运行在事件循环线程上,这个线程两秒内就不能服务其他连接。少量这样的请求足以让延迟迅速扩大。
Reactor 本身不强制每个操作符切换线程。大多数操作符会继续在传来信号的线程上工作。map 不会自动开新线程,flatMap 也不等于并行;内部 Publisher 是否异步,要看它自己的实现。
调度器提供了显式切换执行上下文的能力:
subscribeOn 影响订阅源和向上游请求发生在哪个调度器上,通常紧跟在需要隔离的源后面。publishOn 影响它之后的下游操作符在哪个调度器上处理信号,位置会改变语义。Schedulers.boundedElastic() 为不得不保留的阻塞工作提供有上限的弹性线程资源。Schedulers.parallel() 面向 CPU 计算,线程数通常与处理器核心数相关,它不是阻塞 I/O 的等待室。更深入的线程切换、并发数量和背压关系留到第十二章。现阶段先记住一句能救命的话:事件循环上的每一段用户代码都应尽快返回。
可以临时把线程名放进日志,观察而不是猜测:
return taskCatalog.findAll()
.doOnSubscribe(ignored -> log.info(
"订阅线程:{}", Thread.currentThread().getName()))
.doOnNext(task -> log.info(
"处理任务 {},线程:{}",
task.id(),
Thread.currentThread().getName()));在 Reactor Netty 服务里,常会看到类似 reactor-http-nio-2 的名字:
订阅线程:reactor-http-nio-2
处理任务 1,线程:reactor-http-nio-2
处理任务 2,线程:reactor-http-nio-2这并不证明整个应用永远只用一个线程,也不意味着所有订阅都会落在 -2。它说明当前这段没有主动调度切换,信号沿事件循环继续处理。WebClient 若同样使用 Reactor Netty,客户端与服务器默认还可能共享事件循环资源,所以在远程响应回调里做长时间计算,照样会影响网络处理。
如果确实有一段较重但可拆分的 CPU 计算,可以在明确位置 publishOn(Schedulers.parallel()),让后续计算离开事件循环;完成后不必为了“回到原线程”手工切换,WebFlux 能接收来自不同线程的信号。可这仍需要基准测试。序列化、线程切换和任务调度本身都有成本,小数据上的一次简单 map 放到并行调度器通常只会更慢。
不要在 WebFlux 请求链里调用 block()、blockFirst() 或 blockLast()。它们会把异步结果重新变成“当前线程原地等待”,既破坏并发模型,也可能在非阻塞线程上直接触发异常。控制器应返回 Mono 或 Flux,把订阅和响应写入交给框架。
TaskHub 的任务快照只记录 assigneeId,展示卡片时还要向成员目录服务查询中文姓名。阻塞式客户端常见的写法是发出请求,然后让当前线程等响应;WebClient 则返回 Mono 或 Flux,让远程响应成为管道中的后续信号。
先创建一个配置明确的客户端 Bean:
package com.welearn.taskhub.reactive.member;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.reactive.function.client.WebClient;
@Configuration
public class MemberClientConfig {
@Bean
WebClient memberWebClient(
WebClient.Builder builder,
@Value("${clients.member.base-url}") String baseUrl) {
return builder.baseUrl
clients.member.base-url=http://localhost:8090成员响应只保留这一章用得到的字段:
public record MemberProfile(long id, String displayName) {
}接着封装远程调用:
package com.welearn.taskhub.reactive.member;
import org.springframework.http.HttpStatusCode;
import org.springframework.stereotype.Component;
import org.springframework.web.reactive.function.client.WebClient;
import reactor.core.publisher.Mono;
import java.time.Duration;
import java.util.concurrent.TimeoutException;
@Component
public class MemberClient {
private final WebClient webClient;
public MemberClient(WebClient memberWebClient) {
this.webClient = memberWebClient;
}
调用 findMember(101) 时,方法不会立刻拿到 MemberProfile,也不会在这里阻塞到响应回来。它返回一条说明:订阅后发 GET 请求;404 转成“成员不存在”的错误;服务端错误转成“成员服务不可用”;成功响应解码为 MemberProfile;超过 800 毫秒还没得到信号,就以超时错误结束。
现在把任务和成员组合成展示卡片:
public record TaskCard(
long id,
String title,
String status,
int priority,
String assigneeName) {
public static TaskCard from(TaskSnapshot task, MemberProfile member) {
return new TaskCard(
task.id(),
task.title(),
task.status(),
task.priority(),
member.
package com.welearn.taskhub.reactive.task;
import com.welearn.taskhub.reactive.member.MemberClient;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
@Service
public class ReactiveTaskService {
private final TaskCatalog taskCatalog;
private final MemberClient memberClient;
public ReactiveTaskService(TaskCatalog taskCatalog, MemberClient memberClient) {
this.taskCatalog = taskCatalog;
this
flatMapSequential(..., 8) 最多同时推进 8 个成员查询,并按原任务顺序向下游发卡片。并发上限不是越大越好。它应与上游容量、连接池和延迟目标一起测量;不设边界地向一个慢服务发请求,只会把压力推给对方。
flatMap 适合后一步依赖前一步结果:先拿到任务,才能知道 assigneeId。若两次调用彼此独立,可以用 Mono.zip 同时订阅,再在两边都成功后组合结果。比如任务看板摘要同时需要本地计数和告警服务状态:
public Mono<TaskBoardSummary> loadSummary() {
Mono<Long> unfinished = taskRepository.countByStatus("TODO");
Mono<Integer> activeAlerts = alertClient.countActiveAlerts();
return Mono.zip(unfinished, activeAlerts)
.map(tuple -> new TaskBoardSummary(
tuple.getT1(),
tuple.getT2()));
若两个上游各需要 300 毫秒并且能并发推进,整体等待可能接近较慢的那一个,而不是简单相加为 600 毫秒。这里的“同时”来自两个异步 Publisher 都被订阅,不是 zip 为每个调用新建线程。
zip 也有明确的失败语义:任一 Mono 出错,组合结果就出错,另一个可能被取消;任一 Mono 为空,组合也无法产生那一个最终值。若告警服务是可选信息,应在 zip 之前为它定义精准降级:
Mono<Integer> activeAlerts = alertClient.countActiveAlerts()
.onErrorResume(
AlertServiceUnavailableException.class,
error -> Mono.just(0));是否用 0 代表“不可用”要看 API 契约。很多时候返回 alertStatus: "UNAVAILABLE" 比伪造零条告警更诚实。响应式操作符能写出恢复链,但数据含义仍由我们负责。

控制器依旧很薄:
@RestController
@RequestMapping("/api/reactive-tasks")
public class ReactiveTaskController {
private final ReactiveTaskService taskService;
public ReactiveTaskController(ReactiveTaskService taskService) {
this.taskService = taskService;
}
@GetMapping
public Flux<TaskCard> list() {
return taskService.listCards();
}
启动成员目录和响应式应用后,请求列表:
curl -s http://localhost:8081/api/reactive-tasks有限的 Flux<TaskCard> 会编码成 JSON 数组:
[
{
"id": 1,
"title": "补上接口测试",
"status": "TODO",
"priority": 5,
"assigneeName": "林晓"
},
{
"id": 2,
"title": "整理部署清单",
"status": "IN_PROGRESS",
"priority": 4,
"assigneeName"
这里最重要的并不是链式写法变短了,而是等待成员服务期间,事件循环线程可以处理其他连接。若成员服务自己仍然用阻塞客户端,这个服务端请求链是否端到端非阻塞还要继续向下检查,不能只在 TaskHub 这一端宣布胜利。
WebClient 还有两个工程边界需要提前说明。第一,retrieve() 适合“根据状态码决定成功或错误,然后解码响应体”的常见路径;需要同时读取状态、响应头和不同类型响应体时,可以选择更底层的交换 API,但必须确保响应体被消费或释放。第二,超时应尽量贴近调用边界。把一个很大的总超时放在控制器最外层,只能告诉你整条链慢了,不能说明是连接、读取还是哪一个上游超时。
不要在 WebFlux Service 中这样使用 WebClient:
MemberProfile member = webClient.get()
.uri("/api/members/{id}", id)
.retrieve()
.bodyToMono(MemberProfile.class)
.block();这段代码虽然用了 WebClient,却在最后把结果阻塞取出,前面建立的非阻塞链到这里就断了。正确做法不是寻找另一个“不会报错的 block”,而是让方法返回 Mono<MemberProfile>,调用者继续用 map、flatMap 或 zip 组合,直到控制器把 Publisher 交给框架。
现在回到这一章的主任务:GET http://localhost:8081/api/task-events 建立一条 Server-Sent Events 连接。SSE 使用普通 HTTP,服务器可以沿一条连接不断向浏览器发送文本事件。它适合“服务器单向推送,客户端只需接收”的任务动态;如果双方都要高频双向通信,再考虑 WebSocket。
我们先用一个有限事件源把协议跑通。事件记录四个字段:序号、任务标题、状态和发生时间。
package com.welearn.taskhub.reactive;
import java.time.Instant;
public record TaskEvent(
long sequence,
String taskTitle,
String status,
Instant occurredAt) {
}TaskEventService 每隔 250 毫秒发一条事件,共发四条:
package com.welearn.taskhub.reactive;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;
import java.time.Duration;
import java.time.Instant;
import java.util.List;
@Service
public class TaskEventService {
private static final List<String> TITLES = List.of(
"设计任务接口", "接入数据库", "补上测试", "准备部署");
Flux.interval 是时间源,第一个参数为零表示订阅后立即发第一个序号,之后每 250 毫秒再发一个。map 把 0、1、2、3 转成业务事件。take(4) 收到第四条后取消上游计时器并正常完成,这个边界非常重要;如果没有 take,interval 会一直运行。
这还是一条冷流:每个 HTTP 订阅都会得到自己从序号 1 开始的四条事件。它方便我们稳定验证 SSE 编码和事件顺序,还不是真实的全局任务广播。等协议跑通后,再把这个源替换成共享的热事件源。
控制器只负责声明媒体类型,并把 Flux 原样交给 WebFlux:
package com.welearn.taskhub.reactive;
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import reactor.core.publisher.Flux;
@RestController
@RequestMapping("/api/task-events")
public class TaskEventController {
private final TaskEventService service;
public TaskEventController(TaskEventService service) {
this.service =
produces = MediaType.TEXT_EVENT_STREAM_VALUE 告诉 WebFlux 不要等 Flux 完成后再拼一个普通 JSON 数组。每个 TaskEvent 到达时,编码器就写出一个 data: 帧并刷新。框架订阅 Service 返回的 Flux,因此控制器里仍然不需要 subscribe()。
打开一个终端,使用 -N 关闭 curl 的输出缓冲:
curl -N http://localhost:8081/api/task-events这条命令实际收到四条连续事件,并在最后一条 DONE 后结束:
data:{"sequence":1,"taskTitle":"设计任务接口","status":"IN_PROGRESS","occurredAt":"2026-08-17T08:46:54.018138Z"}
data:{"sequence":2,"taskTitle":"接入数据库","status":"IN_PROGRESS","occurredAt":"2026-08-17T08:46:54.273018Z"}
data:{"sequence":3,"taskTitle":"补上测试","status":"IN_PROGRESS","occurredAt":"2026-08-17T08:46:54.518590Z"}
data:{"sequence":4,"taskTitle":"准备部署","status":"DONE","occurredAt":"2026-08-17T08:46:54.773036Z"}时间来自事件真正创建的瞬间,所以每次运行会不同;稳定的是序号、标题顺序和最后状态。空行是 SSE 事件边界,不要删除。
真实任务变更不会因某个浏览器刚好连接才开始。我们可以用 Reactor Sinks 接住应用内发生的任务事件,再把它暴露为热 Flux:
@Component
public class LiveTaskEventFeed {
private final Sinks.Many<TaskEvent> sink =
Sinks.many().multicast().directBestEffort();
public Flux<TaskEvent> stream() {
return sink.asFlux();
}
public Sinks.EmitResult publish(TaskEvent event) {
return sink.tryEmitNext
multicast() 表示当前多个订阅者共享同一批后续事件;directBestEffort() 不会为了慢订阅者建立一个无限缓冲区。没有人在线时,事件也不会自动保存。生产事件的一端必须检查 EmitResult:tryEmitNext 是一次尝试,不会因为方法名叫 publish 就保证交付。多个线程同时发送、订阅者取消或源已终止,都可能失败。
把控制器的源从 service.events() 换成 liveTaskEventFeed.stream() 后,SSE 就会保持连接并等待后续事件。实际部署时通常还会合并一条心跳流,避免长时间没有业务事件时被代理当成闲置连接。心跳只证明连接仍通,不代表业务事件没有丢失。
浏览器端可以用 EventSource 订阅当前只含 data 的默认事件:
<script>
const events = new EventSource("http://localhost:8081/api/task-events");
events.onmessage = event => {
const taskEvent = JSON.parse(event.data);
console.log(`${taskEvent.taskTitle}:${taskEvent.status}`);
};
events.onerror =
当前控制器返回 Flux<TaskEvent>,所以 WebFlux 只需要为我们生成 data。若还需要事件类型、id、重连间隔或注释心跳,可以让控制器返回 Flux<ServerSentEvent<TaskEvent>>,再用 builder 明确设置字段。一条完整 SSE 帧常见的字段各有用途:
data 是业务数据,多行 data 会由客户端按规则拼接。这里放 JSON。event 是事件类型,浏览器用 addEventListener 按类型监听;省略时走默认 message 事件。id 是事件游标。浏览器重连时可以把最后收到的 id 放进 Last-Event-ID 请求头。retry 可以提示浏览器重连等待时间,但最终重连策略仍应考虑服务端容量和抖动。有了 id 不等于自动断点续传。服务端若没有保存这个 id 之后的历史,收到 Last-Event-ID 也没有内容可补。若需求要求“断网一分钟后不能漏任何任务变更”,就应让事件来自可重放日志,控制器根据最后 id 从存储续读,再切换到实时流。这比单纯把 Sink 改成更大的内存 replay 更可靠。
浏览器页面还要考虑跨域。如果 TaskHub 页面来自 8080,事件服务来自 8081,它们是不同源。开发时可以显式配置只允许页面源访问 SSE 端点,生产中也可以通过反向代理把它们统一到同一站点路径。不要为了省事给所有来源、所有方法开放 CORS,尤其是未来连接带凭据时。
代理缓冲也会影响“实时”体验。应用已经逐帧发送,不代表反向代理一定立刻转发;若代理攒够一批才下发,用户会看到事件成团出现。排查时要从事件产生时间、应用写出时间、代理转发和浏览器接收四个位置对照,而不是只盯着 Controller。
客户端关闭页面或执行 events.close() 时,HTTP 订阅会取消,取消信号沿管道向上游传播。能感知取消并及时停止计时器、释放连接,是响应式资源管理的一部分。第十二章会用 doFinally 观察 complete、error 和 cancel 三种收尾方式。
页面收到事件后,也不一定要把整个任务列表重新请求一次。事件可以带 taskId、新状态和版本号,前端只更新对应行;若发现版本跳跃,再主动刷新列表。这样 SSE 负责“告诉你什么变了”,REST 仍负责“给我当前完整状态”。把通知通道与查询通道分开,断线恢复和权限控制都会更清楚。
进程内 Sink 不是可靠消息系统。应用重启、无人订阅或消费者过慢时,事件可能丢失。如果“任务状态变更必须最终送达”是业务承诺,应使用事务消息表、消息代理或可重放事件存储,并用事件 id 处理断线续传。SSE 只是传输通道,不会替你提供持久化可靠性。
命令式代码里,我们习惯在当前调用栈上 throw,再由外层 try/catch 处理。响应式管道真正执行时,原始方法往往早已返回,所以错误要作为 onError 信号沿着订阅链传递。
下面这个 try/catch 看似保护了远程调用,实际上抓不到订阅后才到来的超时:
try {
return memberClient.findMember(id);
} catch (MemberServiceTimeoutException exception) {
return Mono.empty();
}调用 findMember 时只是装配并返回 Mono,方法没有在这一刻等待 800 毫秒,因此 catch 块通常不会执行。超时发生在后续订阅阶段,应在管道中使用 onErrorResume、onErrorMap 等操作符,或者让错误继续到统一异常处理层。只有组装代码本身同步抛出的异常,才可能被这里的 catch 捕获。
先分清“空”和“错”。taskCatalog.findById(999) 没有元素,是正常完成的空 Mono。只有把它转换后,才得到 404 语义:
public Mono<TaskCard> findCard(long id) {
return taskCatalog.findById(id)
.switchIfEmpty(Mono.error(new ReactiveTaskNotFoundException(id)))
.flatMap(this::enrich);
}switchIfEmpty 只处理“没有值但正常完成”,不会吞掉远程超时。这个区别非常实用:没有这个任务是领域结果,成员目录连接失败则是基础设施故障,两者应映射成不同的 HTTP 状态和 ProblemDetail。
如果成员服务暂时不可用,但任务列表允许降级为“负责人暂不可用”,可以只恢复这一类错误:
private Mono<TaskCard> enrich(TaskSnapshot task) {
return memberClient.findMember(task.assigneeId())
.map(member -> TaskCard.from(task, member))
.onErrorResume(
MemberServiceUnavailableException.class,
error -> Mono.just(new TaskCard(
task.id(),
task.title(),
task.
onErrorResume 并不是让失败的原序列从出错位置接着跑。原上游已经终止,它把错误信号替换成另一条 Publisher。onErrorReturn 是固定值版本;onErrorMap 用来保留失败语义但换成更贴近业务的异常;doOnError 适合记录观察信息,不负责恢复。
假设源依次产生任务 1、任务 2、任务 3,而处理任务 2 时出错,那么信号顺序可能是:
onNext(任务 1)
onError(成员服务不可用)任务 3 不会从原序列继续出现。若 onErrorResume 换成一个降级卡片,后续看到的是“任务 1、降级卡片、完成”,仍然不是回到原源继续处理任务 3。若业务要求单个元素失败不影响其他元素,应把错误边界放进每个元素的内部 Publisher,并明确失败元素是跳过、替换还是进入旁路;这个位置差异会改变整个列表的语义。

不要在最外层随手写:
.onErrorResume(error -> Mono.empty())这会把超时、反序列化失败、程序缺陷全都伪装成“没有数据”。客户端看到的是一个安静的空响应,日志也可能没有线索。恢复策略必须对应明确的异常类型,并留下足够上下文。
重试同样不能凭感觉加。retry(3) 的本质是失败后重新订阅上游。对冷的 WebClient Mono 来说,就是再发请求;对创建任务这类非幂等操作,可能造成重复写入;对“成员 id 不存在”的 404,重试只会徒增流量。超时、退避、按异常筛选重试和幂等边界,会在第十二章单独处理。
WebFlux 也支持集中异常处理,控制器不需要每个方法重复 onErrorResume 拼错误 JSON:
@RestControllerAdvice
public class ReactiveApiExceptionHandler {
@ExceptionHandler(ReactiveTaskNotFoundException.class)
public ResponseEntity<ProblemDetail> handleTaskNotFound(
ReactiveTaskNotFoundException exception) {
ProblemDetail problem = ProblemDetail.forStatus(HttpStatus.NOT_FOUND);
problem.setTitle("任务不存在");
problem.setDetail(exception.getMessage());
return ResponseEntity.status(HttpStatus.NOT_FOUND).
异常处理方法这里直接返回 ResponseEntity 没有问题,因为它只是在内存里构造一个很小的响应对象,并不执行阻塞 I/O。如果错误恢复本身还要异步查询,方法也可以返回 Mono<ResponseEntity<ProblemDetail>>。
对于已经写出若干 SSE 帧的连接,情况不同:HTTP 状态和响应头早已发送,后续错误不能再改成一份全新的 500 JSON。流通常会终止并断开连接,客户端根据策略重连。因此,持续流的错误边界、心跳和恢复协议需要在发出第一帧之前设计清楚。
到目前为止,响应式应用使用内存事件源,不是假装数据库问题不存在,而是有意把线程模型讲清楚。真正需要在响应式请求链里访问关系型数据库时,可以考虑 R2DBC。它的驱动契约围绕延迟执行、Reactive Streams 和非阻塞 I/O 设计,Spring Data R2DBC 则提供响应式 Repository 与映射支持。
依赖大致如下:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-r2dbc</artifactId>
</dependency>
<dependency>
<groupId>org.postgresql</groupId>
<artifactId>r2dbc-postgresql</artifactId>
<scope>runtime</scope>
</dependency>连接配置也使用 spring.r2dbc,URL 以 r2dbc: 开头:
spring.r2dbc.url=r2dbc:postgresql://localhost:5432/taskhub
spring.r2dbc.username=taskhub
spring.r2dbc.password=taskhub一个最小的响应式任务映射和 Repository 可以这样写:
package com.welearn.taskhub.reactive.task;
import org.springframework.data.annotation.Id;
import org.springframework.data.relational.core.mapping.Column;
import org.springframework.data.relational.core.mapping.Table;
@Table("tasks")
public record ReactiveTask(
@Id Long id,
String title,
String status,
int priority,
@Column("assignee_id") Long assigneeId) {
}package com.welearn.taskhub.reactive.task;
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
import reactor.core.publisher.Flux;
public interface ReactiveTaskRepository
extends ReactiveCrudRepository<ReactiveTask, Long> {
Flux<ReactiveTask> findAllByStatusOrderByPriorityDesc(String status);
}findById 返回 Mono<ReactiveTask>,findAll 和派生查询返回 Flux<ReactiveTask>。数据库驱动在连接可读写时推动结果信号,不需要用一个事件循环线程停在那里等网络响应。

不过,R2DBC 不是“把 JpaRepository 的返回类型换成 Mono”这么简单。JPA 的持久化上下文、脏检查、延迟加载和实体关系管理有自己的一套模型;Spring Data R2DBC 更接近显式的聚合映射和响应式 SQL 访问。迁移前要检查事务、级联、复杂查询和驱动支持,不能只做接口替换。
事务也要留在响应式链内部。响应式事务管理器把事务状态与订阅上下文关联,而不是依赖某个请求从头到尾固定在线程局部变量里。被事务代理的方法应返回 Publisher,让订阅、数据库操作和终止信号都发生在代理能够管理的范围内。若方法内部偷偷 subscribe(),那条支线可能已经离开事务边界;若先 block() 取结果再返回,响应式事务也失去了意义。
关系型数据库的表结构迁移仍然要做。R2DBC 负责运行期非阻塞访问,不代表自动替你管理版本化 DDL。实际项目通常仍使用 Flyway 或 Liquibase 维护迁移;迁移可以在启动阶段通过独立连接完成,重点是不要把阻塞迁移操作放进每个 WebFlux 请求。
这也是我们没有让主 TaskHub 改用 R2DBC 的原因。它当前围绕 JPA 建立的事务与实体模型工作良好,主要需求是普通 CRUD。为了一个 SSE 端点重写整条数据访问链,代价并不合理。
现实项目里常有一段时间无法马上替换旧 SDK 或 JDBC 调用。过渡期可以把阻塞源放到 boundedElastic:
public Mono<TaskSnapshot> loadLegacyTask(long id) {
return Mono.fromCallable(() -> legacyTaskGateway.findById(id))
.subscribeOn(Schedulers.boundedElastic())
.flatMap(optional -> Mono.justOrEmpty(optional));
}Mono.fromCallable 把调用推迟到订阅时,subscribeOn 让源在 bounded elastic worker 上运行。事件循环不会被这次等待卡住,但阻塞并没有消失:仍然有一个线程在等,队列和线程数量也有上限。隔离可以保护事件循环,是迁移的缓冲带,不是把阻塞库“响应式化”的魔法。
如果大多数请求最终都要这样包住 JPA、文件读取和旧 HTTP 客户端,说明 MVC 更符合应用现实。把整套阻塞程序藏在 boundedElastic 后面,只会同时承担两种编程模型的复杂度。
桥接代码还要设置容量意识。boundedElastic 有线程和排队上限,流量突增时任务仍会等待。如果旧网关有自己的连接池,调度器并发也不应远高于连接池容量,否则大量任务只是从事件循环转移到另一个队列。监控时要同时看 bounded elastic 排队、旧资源池等待和请求超时,不能把“事件循环没被卡”当成整个系统健康。
现在可以把 TaskHub 的响应式路径连起来看一遍。客户端请求一个任务卡片时,过程是:
GET /api/reactive-tasks/1,WebFlux 找到控制器方法。Mono<TaskCard> 并返回。TaskCatalog.findById(1) 发出任务快照;如果为空,switchIfEmpty 发出任务不存在错误。flatMap 根据 assigneeId 订阅 WebClient 返回的成员 Mono。map 组装 TaskCard,WebFlux 把它编码为 JSON。onComplete,HTTP 响应完成;若客户端提前断开,则订阅取消。这里没有一个“Reactor 总调度线程”替你按顺序执行所有代码。操作符围绕订阅关系传递信号,网络运行时在 I/O 就绪时继续推动管道。理解这一点后,链式代码就不再像某种难懂语法,而是一张可执行的数据流图。
对于 /api/task-events,前四步类似,但后面变成一个长期订阅。每来一个任务事件,就写一帧 SSE;没有事件时连接仍在;页面关闭时取消。这也解释了为什么不能把所有 Flux 都 collectList():对无限流来说,等待“收集完”意味着永远没有响应,而且内存会持续增长。
当你评审一段响应式代码时,可以顺着这条路径问五个问题:谁是数据源,谁最终订阅,哪一步可能阻塞,空值与错误分别怎样表示,客户端取消后资源怎样释放。只要其中一个问题答不上来,就先不要急着加更多操作符。复杂链条最怕的不是代码长,而是生命周期没有主人。
如果要决定一个新接口放进主应用还是响应式应用,可以再走一遍具体判断。先看输出:一次性 JSON 还是可能持续很久的流;再看输入依赖:JPA/JDBC 还是 WebClient/R2DBC;然后看并发特征:连接是在做计算,还是大部分时间等网络;最后看失败策略:能否取消上游、是否需要事件重放、调用能不能安全重试。TaskHub 的任务详情接口四项都偏向 MVC,任务动态则明显偏向 WebFlux。
也不要忽略维护成本。命令式代码的调用栈更接近源码顺序,断点调试和线程局部上下文也更直观;响应式代码把控制流交给订阅链,换来组合异步步骤和传播取消的能力。团队若只会复制操作符,很容易写出能编译却无法解释的管道。一个稳妥做法是先挑一条边界清楚、收益可测的流式接口,记录线程数、内存、吞吐与尾延迟,再决定是否扩展,而不是先制定“全面响应式改造”目标。
本章的双应用安排正是这个小步策略。它允许任务动态独立承压,也让 MVC 主应用继续按熟悉的事务模型演进。两个应用之间若将来通过消息代理传递任务变更,事件契约会成为新的边界;若流量很小,也可以把通知能力收回主系统。架构可以随着证据调整,不需要用框架选择证明项目先进。
最后还要观察最慢的那一段。入口使用 WebFlux、出口使用 WebClient,并不代表中间转换没有风险。一次无界 JSON 聚合、一个意外的同步 DNS 查询、一段在事件循环里执行的压缩,都可能重新引入长停顿。端到端非阻塞是一条链的性质,需要用指标和线程采样持续验证,不能靠代码表面全是 Mono、Flux 就下结论。
肉眼看见 curl 打出两条 SSE,只能说明某一次运行碰巧符合预期。Reactor 的延迟执行、错误和取消都需要真正订阅后才出现,因此测试也要用订阅者的视角描述信号。
reactor-test 提供的 StepVerifier 可以逐步声明期望。事件源包含时间操作符,测试用虚拟时间推进一秒,不让构建真的依赖定时等待:
@Test
void emitsFourOrderedEvents() {
StepVerifier.withVirtualTime(service::events)
.thenAwait(Duration.ofSeconds(1))
.expectNextMatches(event ->
event.sequence() == 1
&& event.taskTitle().equals("设计任务接口"))
.expectNextCount(2)
.
withVirtualTime 接收 Supplier,确保 events() 在虚拟调度器安装后才创建。thenAwait 推进时钟,四条 interval 事件立即可验证。测试检查第一条标题、略过中间两条、检查第四条状态,再断言完成信号。若忘了 .take(4),最后的 verifyComplete() 就无法通过。
错误也应按信号断言:
@Test
void missingTaskEndsWithNotFoundError() {
StepVerifier.create(taskService.findCard(999))
.expectError(ReactiveTaskNotFoundException.class)
.verify();
}持续 SSE 不会自然完成,测试必须主动取消,不能让测试进程一直等:
@Test
void eventFeedEmitsPublishedEvent() {
TaskEvent event = new TaskEvent(
1,
"补上测试",
"IN_PROGRESS",
Instant.parse("2026-08-17T09:10:00Z"));
StepVerifier.create(liveTaskEventFeed.stream())
.then(() -> liveTaskEventFeed.publish(event))
.expectNext(event)
.thenCancel
运行响应式项目的完整测试:
./mvnw test结果中四项测试全部通过:
[INFO] Tests run: 4, Failures: 0, Errors: 0, Skipped: 0
[INFO] BUILD SUCCESS这里的四项包括应用上下文、四条有序事件,以及临时故障重试和业务错误不重试。数字本身不是质量证明,但它把本章已经承诺的信号顺序和错误分类固定了下来。
控制器契约还可以用 WebTestClient 验证,而不是在测试里真的启动 curl:
@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
class ReactiveTaskControllerTest {
@Autowired
WebTestClient webTestClient;
@Test
void returnsTaskCardsAsJson() {
webTestClient.get()
.uri("/api/reactive-tasks")
.exchange()
.expectStatus().isOk()
.expectHeader()
.
WebTestClient 自己也沿响应式方式交换请求和响应,能断言状态、响应头和响应体。Service 的转换规则用 StepVerifier 测,Controller 的 HTTP 映射用 WebTestClient 测,职责会比一个巨大的全链路测试更清楚。外部 MemberClient 则应替换为可控的测试服务器,验证 200、404、5xx 和超时映射,避免单元测试依赖真实网络。
时间型操作符尤其不能靠真实等待。若一个超时是 30 秒,让测试每次睡 30 秒既慢又不稳定。Reactor 的虚拟时间可以把时钟推进写进测试,不过操作符必须在虚拟调度器生效后延迟创建。我们先记住这个约束,第十二章会完整写出超时和退避测试。
下一章会把测试范围扩展到两条路径:MVC 的 Service、JPA 和 HTTP 契约,以及响应式流的值、顺序、完成和错误。第十二章再回到这里,用虚拟时间测试超时与退避,用可控请求量观察背压,并验证取消时资源是否真的释放。
这一章没有把 TaskHub 的 MVC + JPA 主应用推倒重来。我们为持续任务动态建立了独立的 taskhub-reactive 应用,并把两条线程模型的边界保留下来:普通 CRUD 继续走容易理解的阻塞式事务路径,SSE 和非阻塞外部调用走 WebFlux。
现在你应该能准确说出 Mono 和 Flux 的含义:它们描述 0..1 和 0..N 的信号序列,不代表后台线程,也不保证并行。操作链先装配,订阅后才执行;WebFlux 控制器返回 Publisher,由框架负责订阅。错误是终止信号,恢复操作符建立的是替代序列。冷流为每次订阅重新开始,热流则让订阅者共享正在发生的事件。
更重要的是,我们给“非阻塞”划了边界。WebClient 和 R2DBC 可以参与非阻塞链,JPA、JDBC、Thread.sleep 和 block() 不会因为外面套了 Mono 就改变性质。boundedElastic 能保护事件循环,但只是过渡隔离。响应式的主要收益是在 I/O 延迟和大量并发连接下,用更少线程获得更可预测的资源占用,而不是让每个请求跑得更快。
下一章先把 TaskHub 的同步与响应式路径都放进自动化测试。等这些基础行为有了可靠断言,第十二章再深入背压、并发上限、超时、重试、取消和调试;那些内容如果没有测试兜底,很容易只停留在“看起来能跑”。