Flux 与 Mono 是 Project Reactor 的两种核心类型:Mono 表示 0 或 1 个元素,Flux 表示 0 到多个元素。学会它们就是学会把“按步骤做事”的阻塞代码,改成“描述数据流、响应事件”的非阻塞风格:创建流、变换、组合、处理错误与背压,再决定在哪个线程上执行,是整个流程的脉络。

先把概念讲清楚:为什么需要 Flux 和 Mono
想象你在排队买咖啡。传统代码像是你自己跑到柜台、点单、等咖啡做好再回到座位:一步步阻塞等待。响应式编程更像你留个号码牌,工作人员准备好了就叫你名字——你不必一直盯着。Flux/Mono 就是那套排队和通知的语言。
Mono 和 Flux 的区别
| 类型 | 语义 | 典型场景 |
| Mono<T> | 0 或 1 个元素(可能是 empty 或 error) | 异步获取单个对象、登录返回 token、单次数据库查询 |
| Flux<T> | 0 到 N 个元素(可能是无限流) | 文件行流、消息队列消费、数据批处理 |
如何创建 Flux / Mono(常见工厂方法)
创建其实不难,先列几个常用的工厂方法,配点说明让它像实物一样容易抓住。
- Mono.just(value):立刻包含一个值。
- Mono.empty():表示无值完成。
- Mono.error(ex):立即终止并抛出错误。
- Flux.fromIterable(list):把集合变成流。
- Flux.range(start, count):生成一串整数。
- Flux.interval(Duration):按时间间隔发射项(常用于模拟周期事件)。
示例代码(伪码,方便理解):
// Mono
Mono m = Mono.just("hello");
Mono empty = Mono.empty();
Mono err = Mono.error(new RuntimeException("fail"));
// Flux
Flux f = Flux.fromIterable(Arrays.asList(1,2,3));
Flux ticks = Flux.interval(Duration.ofSeconds(1));
常用操作符:像拼积木一样组合变换
把操作符想成厨房里的刀具:map 切片,flatMap 是把两个菜合并煮,filter 是挑选好的食材。常见的操作:
- map:同步转换每个元素。
- flatMap:把元素映射成另一个异步流并合并结果(注意并发行为)。
- concatMap:类似 flatMap,但保持顺序串行合并。
- filter:过滤元素。
- take、skip:截断或跳过元素。
- zip:把多个流的相同索引元素合并成元组/对象。
- merge:并行合并多个流,顺序不保证。
比如:从用户 ID 列表并发拉取详情,但保留原有顺序可以用 concatMap,想更快用 flatMap(但要注意并发限制)。
错误处理与重试策略
错误处理在响应式中要显式写出来,否则流会直接终止。常用方法:
- onErrorReturn(value):出错时返回默认值。
- onErrorResume(e -> Mono.just(…)):根据错误动态切换备用流。
- retry(n) / retryWhen:重试策略,可配合延迟和限次。
示例思路:调用远程服务超时,用 onErrorResume 切到缓存值或降级逻辑,retryWhen 则可以在网络短暂抖动时尝试多次。
背压(Backpressure)和调度器(Schedulers)
背压是核心概念:消费者处理不过来,生产者不能无限制发。Flux 支持背压语义,而各种操作符在内部会处理请求多少元素。
Schedulers:控制在哪个线程跑
常用的有:Schedulers.immediate()、Schedulers.boundedElastic()(适合阻塞 I/O)、Schedulers.parallel()(CPU 密集型)。需要知道:
- publishOn:切换之后的操作在指定线程池执行(影响下游)。
- subscribeOn:决定订阅开始在哪个线程执行(影响上游生产)。
| 操作 | 效果 |
| subscribeOn | 改变流的订阅执行位置(上游) |
| publishOn | 改变流之后操作的线程(下游) |
小心点:在一个链里多次 publishOn 会在链中多次切换线程,理解这点能避免调试噩梦。
组合多个流:常见场景与选择
你可能会把多个数据源合并:数据库、缓存、HTTP。选择合适的组合方式很重要:
- zip:用来将多个单次结果(Mono)按顺序组合成一个复合对象,常用于并行请求多个接口然后合并结果。
- merge:把多个 Flux 的元素交错发出,适合事件流合并。
- concat:串行拼接,按顺序等待前一个完成再发下一个。
- flatMapSequential:并发获取但按源顺序输出(权衡并发与顺序)。
与 Spring WebFlux 的整合要点
在 WebFlux 控制器中直接返回 Mono/Flux:Spring 会把它们转成响应。注意几点:
- 不要在 Reactor 流里做阻塞调用(如 JDBC 的阻塞 I/O),否则需切到 boundedElastic 并注明。
- 数据库要用 R2DBC 或响应式驱动;Redis、Mongo 也有响应式客户端。
- 对于文件上传/下载、SSE(Server-Sent Events)等场景,Flux 非常合适。
测试与调试小技巧
调试响应式链条有时像跟踪流水线,工具和方法很关键:
- StepVerifier:来自 reactor-test,用于断言流的行为和顺序。
- log():在流中插入 .log() 可以看到信号(onSubscribe、request、onNext、onComplete、onError)。
- Block 用在测试环境:不要把 block() 放在生产代码,测试时可短暂用 block() 验证结果。
常见陷阱与最佳实践(经验之谈)
- 别在响应式链里混入大量阻塞调用;如果不得不阻塞,隔离到 boundedElastic。
- 理解 flatMap 的并发语义:默认并行且无序,可能导致顺序问题与资源争用。
- 对无限流(如 Flux.interval)要有取消策略,避免泄露订阅。
- 在高并发下,尽量限制并发度(flatMap 的 concurrency 参数、limitRate 等)。
- 良好地处理错误和超时:Web 客户端和 DB 调用都应设置超时和降级方案。
实战小案例:按 ID 并发拉取详情并合并(保持原始顺序)
Flux.just(1,2,3,4)
.concatMap(id -> fetchDetailAsync(id)) // 保持顺序,串行或可用 flatMapSequential 控制并发
.map(this::enrich)
.onErrorResume(e -> fallback())
.subscribe(result -> System.out.println("got: " + result));
如果想并发但仍保序,可以用 flatMapSequential 或者 flatMap + index + sort(后者会有延迟与内存成本)。
细节速查表
| 场景 | 推荐 |
| 单次异步结果 | Mono |
| 多元素或事件流 | Flux |
| 阻塞 I/O | boundedElastic + 明确隔离 |
| 保证顺序 | concatMap / flatMapSequential |
写到这里我在想,很多人初学时最困惑的往往不是某个操作符的名字,而是思路——把“你要做什么”和“什么时候做”分清楚:先描述数据流(是什么),再选用操作符(怎么变),最后安排执行环境(在哪儿跑)。照着这个顺序来,Flux/Mono 的世界其实没那么可怕。