Reactor 操作符详解
2026/7/16大约 3 分钟
Reactor 操作符详解
一句话概括
Reactor 的操作符是构建响应式流水线的"乐高积木"——每个操作符包装上一个 Publisher,通过组合串联实现数据的转换、过滤、合并与条件选择。
1. 转换操作符
| 操作符 | 作用 | 顺序保证 | 并发度 |
|---|---|---|---|
map | 同步 1:1 转换 | ✅ 保持顺序 | 串行 |
flatMap | 异步 1:N 展平 | ❌ 无顺序保证 | 并发 |
concatMap | 顺序 1:N 展平 | ✅ 保持顺序 | 串行 |
filter | 按条件过滤 | ✅ 保持顺序 | 串行 |
map — 最基础的数据转换
每个元素同步转换,返回一个值:
Flux.just(1, 2, 3)
.map(i -> "Number: " + i)
// → "Number: 1", "Number: 2", "Number: 3"
// 实际场景:从 Entity 到 DTO 的映射
Flux<UserEntity> users = userRepository.findAll();
users.map(UserDTO::fromEntity); // 批量转换flatMap — 异步并发展平
每个元素展开为一个子流,子流可以异步执行,结果按完成时间到达(可能乱序):
Flux.just(1, 2, 3)
.flatMap(i -> Mono.just("Item-" + i)
.subscribeOn(Schedulers.parallel()));
// 实际场景:根据 ID 并发查询远程 API
Flux<String> ids = Flux.just("a", "b", "c");
Flux<User> users = ids.flatMap(id -> userService.findById(id));
// 三个请求并发,先完成的先返回concatMap — 顺序保序版 flatMap
保持每个子流的顺序,但牺牲并发度——前一个子流完成才启动下一个:
Flux.just(1, 2, 3)
.concatMap(i -> Mono.just("Ordered-" + i)
.delayElement(Duration.ofMillis(50)));
// 即使每个元素延迟 50ms,输出依然按 1, 2, 3 顺序filter — 保留符合条件的元素
Flux.range(1, 10)
.filter(i -> i % 2 == 0)
// → 2, 4, 6, 8, 10执行流程对比
map: 1 ──→ "1" 2 ──→ "2" 3 ──→ "3" (串行,顺序保持)
flatMap: 1 ──→ "1-1" ──→ "1-2"
2 ──→ "2-1" (并发,可能交错)
3 ──→ "3-1" ──→ "3-2"
concatMap: 1 ──→ "1-1" ──→ "1-2" (顺序保序)
2 ──→ "2-1" ──→ "2-2" (等 1 完成再开始 2)2. 合并操作符
| 操作符 | 行为 | 场景 |
|---|---|---|
merge | 交错合并,顺序取决于完成时间 | 实时数据合并 |
concat | 串行合并,前一个完成后才订阅下一个 | 有序数据拼接 |
zip | 一对一配对,所有源都有数据才发射 | 聚合多个 API |
merge — 交错合并
Flux<String> fast = Flux.just("A1", "A2")
.delayElements(Duration.ofMillis(50));
Flux<String> slow = Flux.just("B1", "B2")
.delayElements(Duration.ofMillis(80));
Flux.merge(fast, slow);
// → A1, B1, A2, B2(取决于实际延迟,可能交错)concat — 串行合并
Flux.concat(
Flux.just("A1", "A2"),
Flux.just("B1", "B2")
);
// → A1, A2, B1, B2(严格顺序,B 等 A 完成)zip — 配对合并
Flux.zip(
Flux.just("Alice", "Bob"),
Flux.just(95, 88),
(name, score) -> name + ": " + score
);
// → "Alice: 95", "Bob: 88"
// 实际场景:并行查询用户信息和订单信息,配对组装
Flux.zip(userService.findById(id), orderService.findByUserId(id));3. 操作符选型决策
选择哪个操作符,取决于你的需求:
是否需要保持输入顺序?
├─ 是 → 是否需要异步展开子流?
│ ├─ 是 → concatMap
│ └─ 否 → map(1:1)或 filter(N:1)
└─ 否 → 是否需要异步展开子流?
├─ 是 → flatMap
└─ 否 → merge(合并流)或 zip(配对)经验法则:
| 场景 | 推荐操作符 |
|---|---|
| 简单 1:1 映射(DTO 转换、字段提取) | map |
| 按 ID 并发调用远程服务 | flatMap |
| 保持请求顺序的远程调用 | concatMap |
| 过滤不需要的数据 | filter |
| 合并多个数据源(日志、事件流) | merge |
| 按顺序拼接数据 | concat |
| 聚合多路数据(如 zip 多个 API 结果) | zip |
关键总结
| 对比 | map | flatMap | concatMap |
|---|---|---|---|
| 输出 | 1:1 转换 | 1:N 展平 | 1:N 展平 |
| 顺序 | ✅ 保序 | ❌ 按完成时间 | ✅ 保序 |
| 并发 | 串行 | 并发 | 串行 |
| 场景 | DTO 映射 | 并发远程调用 | 顺序远程调用 |
上一篇:Reactor 核心概念
下一篇:Reactor 错误处理与背压