Reactor 核心概念
2026/7/16大约 3 分钟
Reactor 核心概念
一句话概括
Project Reactor 是 Spring WebFlux 底层的响应式编程库,基于 Reactive Streams 规范提供了 Mono<T>(0-1 个元素)和 Flux<T>(0-N 个元素)两种核心类型,让开发者用声明式 DSL 编排异步、非阻塞的数据流。
1. Reactive Streams 规范
Reactive Streams 定义了 4 个核心接口,是整个响应式生态的基石:
| 接口 | 角色 | 说明 |
|---|---|---|
Publisher<T> | 发布者 | 可订阅的数据源,唯一方法 subscribe(Subscriber) |
Subscriber<T> | 订阅者 | 消费数据的接收者,4 个回调方法 |
Subscription | 订阅契约 | 连接发布者与订阅者,控制背压 request(n) |
Processor<T,R> | 处理器 | 既是发布者又是订阅者 |
背压契约
// 订阅者通过 Subscription 向上游按需请求数据
public void onSubscribe(Subscription s) {
s.request(1); // 请求 1 个元素
}
public void onNext(T item) {
process(item); // 处理元素
subscription.request(1);// 处理完再请求下一个
}Project Reactor 完整实现了这套规范,并在此基础上提供了丰富的操作符 DSL,让开发者可以像组合流水线一样编排异步任务。
2. Mono 与 Flux
| 类型 | 语义 | 类比 |
|---|---|---|
Mono<T> | 0 或 1 个元素的异步序列 | 类似 CompletableFuture<T> + 操作符 |
Flux<T> | 0 到 N 个元素的异步序列 | 类似 Stream<T> + 异步 + 背压 |
创建方式
// ─── Flux ───
Flux<String> f1 = Flux.just("a", "b", "c"); // 从数组
Flux<Integer> f2 = Flux.fromIterable(List.of(1, 2, 3)); // 从集合
Flux<Long> f3 = Flux.range(1, 5).map(i -> (long) i); // 范围
Flux<Long> f4 = Flux.interval(Duration.ofMillis(100)); // 定时发射
Flux<Integer> f5 = Flux.empty(); // 空流
Flux<Integer> f6 = Flux.error(new RuntimeException()); // 错误流
// ─── Mono ───
Mono<String> m1 = Mono.just("hello"); // 有值
Mono<String> m2 = Mono.empty(); // 空
Mono<String> m3 = Mono.fromCallable(() -> "computed"); // 延迟计算冷流 vs 热流
| 特性 | 冷流(Cold) | 热流(Hot) |
|---|---|---|
| 行为 | 每次 subscribe() 重新发射 | 订阅者在任意时刻加入,只收后续数据 |
| 举例 | Flux.just()、Flux.range() | Sinks.many()、Flux.interval() |
| 场景 | 数据计算、数据库查询 | 实时行情、事件总线 |
3. 订阅与执行时机
响应式流是懒加载的——只有调用了 subscribe(),数据才会开始流动:
Flux<Integer> flux = Flux.range(1, 5)
.map(i -> i * 2); // 此时只是定义流水线,什么都没发生
flux.subscribe(i -> System.out.println(i)); // 订阅后才开始执行这种"定义与执行分离"的机制,让 Reactor 可以进行操作符组合优化(如操作符融合),也使得异步流水线的构建变得清晰可控。
4. 操作符流水线
Reactor 最强大的设计之一是操作符(Operators)——每个操作符都会创建一个新的 Publisher 包装上一个,形成一条链式流水线:
Flux.range(1, 5) // 数据源
.filter(i -> i%2==0) // 过滤
.map(i -> "num: "+i) // 转换
.subscribe(...) // 触发执行每个阶段互不干扰,可以独立理解、独立测试。
5. 测试:StepVerifier
Reactor 提供了专门的测试工具 StepVerifier,让异步流的验证变得像写断言一样简单:
// 常规验证
StepVerifier.create(flux)
.expectNext("a", "b", "c") // 断言多个元素
.verifyComplete(); // 验证流正常结束
// 虚拟时间加速(测试 interval 等时间敏感操作)
StepVerifier.withVirtualTime(() -> flux)
.thenAwait(Duration.ofSeconds(1))
.expectNext(0L, 1L, 2L)
.verifyComplete();
// 错误验证
StepVerifier.create(flux)
.expectNext(1)
.expectErrorMessage("oops")
.verify();withVirtualTime 是测试定时操作的神器——不需要真的等 10 秒,虚拟时间瞬间推送到指定时刻。
关键总结
| 概念 | 一句话 |
|---|---|
| Reactive Streams | Publisher-Subscriber 通过 Subscription 实现背压控制的异步流规范 |
| Mono | 0 或 1 个元素,类比异步的 Optional |
| Flux | 0 到 N 个元素,类比异步的 Stream + 背压 |
| 冷流 | 每次订阅重新发射(Flux.just) |
| 热流 | 订阅者加入后只收后续数据(Sinks.many) |
| 懒加载 | 不 subscribe 不执行,定义与执行分离 |
| StepVerifier | 响应式测试利器,支持虚拟时间加速 |
下一篇:Reactor 操作符详解 — 深入 map/flatMap/concatMap/filter/merge/concat/zip 等核心操作符。