Java 响应式编程 - Reactor 基础库
学习 Reactor 绝对不能盲从,因为网上有大量旧版本或错误的旧观念。需要我们死磕官方一手资料:
- 官方网站:projectreactor.io
- 核心文档:Reactor Core 官方参考文档、以及其中 操作符的选择
学习路线
两条管道 Flux 与 Mono
在 万物皆为流 的世界里:
- 数据未到,管道先成:你用声明式代码把“流操作”搭建好时,数据可能根本还没产生(例如用户还没发起请求,数据库还没返回)。
- 数据一到,自发响应:当网卡或者数据库产生了一个数据,这个数据就像一滴水注入了水管。它会自动触发沿途的 filter、map 等流操作,一路向下游涌去。
- 顺流而下,完美返回:最终这个数据被流操作加工成了“新数据”,直接冲刷到前端用户的浏览器上。
这就是为什么响应式编程能做到非阻塞高并发——因为线程不再围着 “方法调用” 转,而是变成了 “管道的维护者”。线程只要看到水管里有水(数据产生),就过去顺着流操作推一把,没水的时候立刻去干别的事。Reactor 的世界里,万物皆可被装进管道。它对 Publisher 规范做到了极致的精简,只给了我们两条核心管道:
- Flux<T>:多元素异步管道
- 定义:代表发出 $0$ 到 $N$ 个 元素的异步响应式序列。
- 场景:读文件行、数据库返回的多行记录(ResultSet)、股市实时行情流、聊天室消息。
- 结束标志:发出完最后一个元素后,触发 onComplete 成功结束,或中途报错触发 onError 强行终止。
- Mono<T>:单元素异步管道
- 定义:代表发出 $0$ 或 $1$ 个 元素的特殊响应式序列。
- 场景:根据 ID 查询用户(要么有,要么没有)、发送异步 HTTP 请求(只返回一个 Response 对象)。
- 结束标志:在 Web 开发中,Mono 的使用频率甚至高于 Flux,因为绝大多数传统的 RPC 调用、单条数据增删改查都是单次返回。
几百个操作符
初学者最大的误区,就是试图去背诵 Reactor 的几百个操作符。其实,根据二八定律,生产环境中高频使用的操作符不超过 20 个。我们将它们归为 4 类最核心的 “积木”:
- ① 创建流(管道水源)
- Flux.just(1, 2, 3):用固定元素直接创建。
- Flux.fromIterable(list):把传统的 Java 集合转成响应式流。
- Mono.justOrEmpty(nullableObj):安全包装一个可能为 null 的对象。
- ② 转换与过滤(管道加工)
- map(t -> r):一对一转换。把苹果变成苹果汁。它是同步的。
- flatMap(t -> Publisher):最核心! 一对一或一对多异步转换。如果你的转换操作本身也需要去调别的微服务或查库(返回 Flux/Mono),必须用 flatMap 把内部的管道平铺拆散,融合成一个主流。
- filter(predicate):条件过滤,不满足的直接在管道中丢弃。
- ③ 组合流(多管齐下)
- Flux.zip(fluxA, fluxB):并行神器。像拉链一样,把两个流的元素按顺序一一对齐组合(如同时查用户信息和订单信息,最后合并返回)。
- mergeWith(otherFlux):把两个流的水汇集到一个流里,谁先到先发谁,不保证顺序。
- ④ 异常与后勤(安全阀)
- onErrorReturn(fallbackValue):降级处理。发生错误时,不让系统崩溃,而是发射一个默认值。
- onErrorResume(e -> …):发生错误时,无缝切换到另一条备用管道。
- doOnNext() / doOnError():只看不改的“窥探”操作,一般用来打日志,不影响数据流本身。
其他重点
在深入 Reactor 时,有三个底层概念是你在未来写代码时一定会遇到并卡壳的,这是学习的重中之重。
- 冷流(Cold)与热流(Hot)的区别
- 冷流(默认):不订阅,不生产。每个订阅者进来,都会重新从头播放一遍数据(像看优酷视频)。
- 热流:不管有没有人看,都在源源不断地广播。订阅者进来只能看到当前及之后的数据(像电视直播)。
- 线程切换控制:publishOn vs subscribeOn
- publishOn 改变的是它之后的操作符在哪个线程执行。
- subscribeOn 改变的是整个流最开始源头的数据加载在哪个线程执行。
- 背压溢出策略
- 当下游的处理速度确实跟不上上游时,除了常规的拉取,Reactor 还提供了 onBackpressureDrop(处理不过来就丢弃)和 onBackpressureBuffer(处理不过来先放缓存区)等具体控制手段。
从简单开始
引入依赖
1 | <dependencyManagement> |
如果你当前的系统是标准的 Spring Boot 项目(如使用了 WebFlux),那么你完全不需要手动去指定 reactor-bom 的版本。因为 Spring Boot 已经在它自己的官方 spring-boot-dependencies 中做好了最佳的版本依赖仲裁。你只需要像下面这样直接引入 WebFlux 或者基础的 starter,它会自动给你带上最兼容、最稳定的对应 Reactor 版本:
1 | <dependencies> |
创建管道
在 Reactor 中,代码的第一步永远是“把数据塞进管道”,变成 Flux 或 Mono。
1 | import reactor.core.publisher.Flux; |
此时的重点:如果你直接运行上面的代码,控制台什么都不会输出。 因为你现在只是画好了“自来水管的设计图”。记住那句至理名言:不订阅,什么都不会发生。
订阅通水
要想让水管里的水流出来,必须在管道的最末端安装一个水龙头,也就是调用 subscribe()。
1 | Flux<String> flux = Flux.just("Java", "Python", "Go"); |
初试管道配件
现在水能流过去了,我们要在水管中间加装 “净水器”。这就是操作符。
1 | Flux<Integer> scoreFlux = Flux.just(45, 78, 92, 59, 100); |
异步编排的灵魂 flatMap
在 Project Reactor 的几百个操作符中,flatMap 是最具分水岭意义的核心操作符,也是初学者最容易卡壳的地方。在命令式编程中,如果我们要对一组数据进行 “再去查一次数据库/调一次微服务” 的操作,通常会用 for 循环里套一个 RPC 调用。但在响应式世界里,这种做法是致命的。要想彻底精通 flatMap,我们先从它的兄弟 map 开始对比,用自来水管的哲学来一击看穿。用一句话概括它们的不同:
- map:同步转换。进去一个元素,出来一个普通元素(一对一)。
- flatMap:异步扁平化转换。进去一个元素,出来一个全新的子管道(Flux/Mono),并把所有子管道融合成一条主流。
如果我们错误地使用 map,会发生什么?
1 | // 假设这是一个异步查库的方法,返回一个 Mono 管道 |
因为 getUserNameFromDb 本身返回的是一个管道(Mono),map 只是机械地把元素包进去。结果导致你得到的是一个 “套娃管道” (Flux)。你装水龙头(subscribe)的时候,流出来的不是水(数据),而是一个个套着的水管!你根本拿不到里面的数据。这时候就需要使用到 flatMap!
1 | Flux<Integer> userIdFlux = Flux.just(1, 2, 3); |
底层流转轨迹:
- 分流:当元素 1 流过 flatMap 时,flatMap 调用 getUserNameFromDb(1),产生了一个独立的、异步的子管道 Mono。
- 并发执行:元素 2、3 过来,同样各自产生了自己的异步子管道。这些子管道在底层同时(并行)开始执行异步查库。
- 扁平化融合:flatMap 内部充当了一个 “总集水器”。哪个子管道的数据先从数据库返回,它就把这滴 “水” 捞出来,顺着主管道送给最终的订阅者。
下面看一个简单的例子:
1 | public static void main(String[] args) throws InterruptedException { |
1 | 🔥 [parallel-1]消费者收到 -> [create in main]订单详情: Order_B | 耗时: 1085ms |
- 总耗时只有 1 秒:我们有 3 个订单,每个耗时 1 秒。如果是传统的 for 循环,总共要死等 3 秒。而 flatMap 让这 3 个操作同时并发出发,最后几乎在同一时间一齐流回主管道!
- 乱序特征:仔细看输出,Order_B 比 Order_A 先出来!因为 flatMap 内部是异步触发、谁先到家谁先走的。它不保证顺序。
如果你的业务场景严格要求顺序(比如必须按 A、B、C 顺序返回),Reactor 提供了另一个管道配件叫 concatMap。它的用法和 flatMap 一模一样,但它内部是同步排队的(耗时会变成 3 秒)。
1 | orderIdFlux.concatMap(FlatMapAdvanced::getOrderDetails).subscribe.. |
1 | 🔥 [parallel-1]消费者收到 -> [create in main]订单详情: Order_A | 耗时: 1492ms |
自定义hanle操作
在日常开发中,我们经常遇到这样的骚操作:拿到一个数据 $\rightarrow$ 先判断它合不合规(filter) $\rightarrow$ 如果合规,把它转成另一种对象(map)。如果用常规写法,你需要连着写 .filter(…).map(…)。但如果使用 handle,它会提供给你一个全能的 SynchronousSink(同步发射器),让你在一个方法块里自发决定这个元素是 “保留并转换”、“直接扔掉”还是“拉响警报(报错)”。注意,虽然它很自由,但它和 generate 一样,在单次处理某一个元素时,最多只能调用一次 sink.next(data)(可以不调用,代表过滤掉;调用一次,代表加工传给下游)。
典型的使用场景,如全能数据清洗仓(过滤 + 转换 + 报错)。假设你正在负责某后台用户评论数据清洗。你需要对一串文章 ID 或者是内容流进行处理:
- 如果内容包含“违规恶意政治”,立刻拉响警报(抛出异常,触发降级)。
- 如果内容包含“广告/水贴”,直接无视它(过滤掉)。
- 如果内容合法,自动给它追加“官方审核”的后缀(转换)。
1 | public class HandleDemo { |
1 | 📝 正常内容A [已通过系统合规审核] |
transform [Deferred]
这两者的本质区别是它们对外部变量的“捕获时机”不同,而这个时机差导致了“是否能获取到外部变量最新值”的本质区别。
- transform(Function) —— 组装期捕获(只看应用启动那一刻)。当你的代码执行到 transform 这一行时(通常是 Spring 容器启动、对象初始化、或者管道声明时),Reactor 就会立刻执行里面的 Function。如果 Function 内部去读取了一个外部变量,它在组装期就把这个变量的值给固定下来了。哪怕这个外部变量是个 AtomicInteger 或者普通的成员变量,之后它的值在运行时变了,后续的所有订阅者,看到的依然是组装期捕获的旧值。
- transformDeferred(Function) —— 订阅期捕获(千人千面,看拧开水龙头那一刻)。代码执行到 transformDeferred 时,它在组装期什么都不做(延迟执行)。只有当某个请求过来说 flux.subscribe()(拧开水龙头)时,它才会当场现去执行里面的 Function。这时 Function 才会去读取外部变量。因此它每次都能抓取到外部变量在 “当下这一秒” 的最新状态。
我们来看下面这个例子:
1 | public static void main(String[] args) { |
1 | 订阅者1;v=A |
将上面的 transform 改成 transformDefer,执行结果就会变成:
1 | 订阅者1;v=A |
其实,这两兄弟的诞生,不是为了对流里的数据做具体的加减乘除,而是为了解决响应式开发中的一个终极痛点:代码如何优雅地实现复用(响应式切面重构)。一句话切中它们的本质区别:
- transform 属于“静态组装”:在代码编译/编译后的链条组装期,逻辑就已经写死了,且全局只执行一次。
- transformDeferred 属于“动态加载”:在每一个订阅者真正调用 subscribe() 的订阅期(Subscription Time)才会触发执行,每个订阅者都会单独动态执行一次。
再来看一个使用 transform 重构的示例:
1 | public class TransformDemo { |
总结起来,只要你的重构逻辑需要动态适配外部配置变更、动态读取本地缓存、或者动态追踪线程上下文(如 MDC/TraceId),必须用 transformDeferred。如果只是纯粹的无状态工具链打包,用 transform 即可。
publishOn / subscribeOn
在 Project Reactor 的异步世界里,publishOn 和 subscribeOn 是两个威力巨大、但极其容易让人混淆的“空间传送门”。它们的本质,就是在管道的中间或源头,强行加装 “动力泵(线程池)”,让水流在不同的物理线程之间横跳。
- publishOn(Scheduler):影响下游。当数据经过这个节点后,它后面的所有加工步骤(如 map, filter)都将切换到指定的线程池中执行。
- subscribeOn(Scheduler):影响源头。不管它被写在管道的哪个位置,它都直接作用于最上游的数据源加载,规定最开始的水是从哪个线程池流出来的。
这里的 Scheduler 可以使用 Reactor 官方的 Schedulers 工具类来生成。它封装了各种满足不同底层工业特性的物理线程池。响应式编程的核心就是用正确的线程池干正确的事。Reactor 官方一共为我们准备了 6 种最常用的现成线程池(调度器)。
- Schedulers.parallel() - 固定线程数池。默认创建的物理线程数量,严格等于你服务器的 CPU 核心数(可以通过配置 reactor.schedulers.defaultPoolSize 修改)。专门为 CPU 密集型任务量身打造。因为线程数和核心数一致,保证了每个核心全速运转的同时,最大化减少了线程上下文切换(Context Switch)的损耗。
- Schedulers.boundedElastic() - 动态扩容的弹性池。它是一个专门用来安放 “脏活、累活、慢活” 的救火队。它的线程数是动态的。默认上限是 CPU 核心数 $\times$ 10。更厉害的是,它自带一个容量高达 100,000 的任务排队队列。如果某个线程闲置超过 60 秒,会自动被回收到内存中,非常节省资源。通常用于传统的阻塞式 JDBC 数据库操作、阻塞式远程调用、本地磁盘文件的 I/O 读写等。
- Schedulers.single() - 全局唯一的单线程池。不管你调用多少次,它在后台对应的永远是同一个物理线程,直到调度器被销毁。而且它会对提交进来的任务实行严格的先进先出(FIFO)排队顺序执行。适用于需要严格保证绝对顺序的串行流处理。比如多线程并发修改某个核心内存状态,为了彻底免去写 synchronized 或 Lock 的代码,直接用 publishOn(Schedulers.single()),让所有修改动作在同一个线程里排队执行,天然免疫并发冲突。
- Schedulers.immediate() - 零线程切换的隐形调度器。它不会创建或切换到任何新的物理线程池中。数据流走到这一步时,当前是谁的线程在执行,就立刻在当前线程原地执行。默认的兜底策略。当你写了一个通用的业务组件,要求默认不在任何线程池里折腾,除非用户显式指定时,就可以把 immediate() 作为默认参数传入。
- Schedulers.newBoundedElastic(…) / newParallel(…) - 私有隔离池定制器。前面介绍的 Schedulers.parallel() 和 boundedElastic() 都是全局共享(Shared)的。也就是说,如果你的项目里有 100 个 Controller 都用了全局池,一旦其中一个 Controller 把池子占满了,其余的都会跟着倒霉。而这两个是专门的私有隔离池定制器。
- Schedulers.elastic() - 无上限的野蛮膨胀池,生产禁用。它是 boundedElastic() 的前身。它的恐怖之处在于线程数完全没有上限。只要有新的阻塞任务进来,它就会无脑地去创建新的物理线程。一旦遇到数据库大面积超时,它会在瞬间疯狂创建成千上万个线程,直接导致服务器 OOM 或者 CPU 因为疯狂切换线程而直接卡死瘫痪。在现代 Reactor 开发中,已彻底被 boundedElastic() 取代。
- Schedulers.fromExecutor(..) - 使用自定义的线程池。
回到正题,我们来看一个 publishOn 和 subscribeOn 实际场景:Netty 网络线程收到请求 $\rightarrow$ 扔给 Parallel(计算密集池)做繁重的加密计算 $\rightarrow$ 扔给 BoundedElastic(弹性阻塞池) 去调用传统的阻塞 MySQL JDBC 驱动。
1 | public static void main(String[] args) throws InterruptedException { |
1 | ✏️ [main] 提取ID: user_101 |
- 初始流由 main 线程触发(在 WebFlux 中则是 Netty 线程)。
- 当数据撞击到第一个 publishOn(Schedulers.parallel()) 时,底层的无锁环形队列发挥作用,数据被递交给 parallel-1 线程,随后的加密 map 立即在计算池中运行。
- 当再次撞击到 publishOn(Schedulers.boundedElastic()) 时,空间再次切换。传统阻塞查库被安放在了专门的弹性阻塞池中,完美保护了主计算池和网络池不被卡死。
为什么说 subscribeOn 是“逆流起效”的?我们来看下面这个案例:
1 | public static void main(String[] args) throws InterruptedException { |
1 | [boundedElastic-1] 加工数据 |
- 订阅期(从下往上):调用 subscribe() 时,信号像多米诺骨牌一样从下往上逆流回溯。
- 触发 subscribeOn:往上回溯遇到 subscribeOn 时,它强行接管:“别往上找了,从现在开始,到最顶端数据源加载的这个逆流动作,由我指定的 boundedElastic-1 线程来干!”
- 运行时(从上往下)**:于是,最顶层的数据源在 boundedElastic-1 线程中被加载,并顺流而下,一路带偏了后面的 map。
这样做的好处是,即使这个数据库查询要耗时 5 秒,由于使用了 subscribeOn(Schedulers.boundedElastic()),整个阻塞过程也只会占用 boundedElastic 池子里的线程。而 WebFlux 核心的那几个极其宝贵的、负责接网卡 HTTP 请求的 Netty Worker 线程依然一身轻松,可以继续疯狂接收其他用户的请求。这就是 Reactor 的“幻影移形”魔法,利用空间换时间,用精细的线程池隔离,达成了整体系统的高可用。
1 | public static void main(String[] args) throws InterruptedException { |
1 | [boundedElastic-1] 加工数据1 |
runOn 切换并行加工
在说明 runOn 之前,我们需要先明确一个痛点:在 Reactor 中,普通的 Flux 管道默认是单线程流水线(串行)工作模式。哪怕你用 publishOn(Schedulers.parallel()) 切换到了多核计算池,流里的数据依然是由池子里的某一个线程,排队一个接一个地处理。如果某个中间加工环节非常消耗 CPU(比如大量复杂的加解密、大文本解析),这种串行排队处理就会浪费服务器的多核性能。runOn 操作符,就是为了彻底解决 “多核 CPU 利用率不足、无法真正并行加工” 的问题而诞生的。
- publishOn(单轨制):3 个任务排成一队。publishOn 把整个队伍打包搬到另一个工位(线程B)。线程B依然需要一个接一个(串行)地完成这 3 个任务。
- parallel().runOn()(多轨制):3 个任务过来。parallel() 把它们拆分成 3 个独立的小组,runOn 直接派 3 个不同的线程(线程B1、线程B2、线程B3)同时(并行)开始干活。
要使用 runOn,必须配合 parallel() 操作符。它们两个是不可分割的 “黄金搭档”。
- parallel(n):它的作用是 “分流”。把一个普通的单列纵队 Flux 管道,切分成 n 个并行的轨道(Rails),此时流的类型会从 Flux 变成 ParallelFlux。
- runOn(Scheduler):它的作用是 “给每个轨道分派物理线程”。它让这 n 个轨道同时在线程池的多个独立线程上并发跑起来。
典型生产案例:多核心并行处理 CPU 密集型任务。假设我们要对一批用户密码进行极其消耗 CPU 的 BCrypt 强哈希加密。如果串行加密,总耗时会随着数据量暴增。我们通过 runOn 压榨多核 CPU 性能。
1 | import reactor.core.publisher.Flux; |
1 | 🔑 HASHED_[pwd1]_BY_parallel-1 |
- 耗时被极度压缩:总共 8 个密码,每个耗时 100ms。如果是传统的串行,总耗时需要 800ms。由于我的机器调度器分成了 4 个并行轨道,4 个物理线程(parallel-1 到 parallel-4)同时开工,每人分摊 2 个任务,最终仅用 200ms 出头就全部搞定了!
- sequential() 的必要性:多轨运行虽然快,但无法直接交给普通的 Web 框架(如 Spring 响应式返回)。必须在加工完后调用 .sequential(),把并行的多轨“合流”重新变成一个普通的 Flux,以便后续的统一投递。
这里有个经典的问题,也就是既然 flatMap 也能并发,为什么还要用 runOn?这里给出它们的本质区别:

也就是如果你写代码是为了不让线程在网络/磁盘 I/O 前死等,用 flatMap。如果你写代码是为了榨干服务器 8 核、16 核 CPU 的算力来做重度计算,用 parallel().runOn()。
冷热流及其切换
理解 “冷流”(Cold Stream)与“热流”(Hot Stream),是迈向 Reactor 高级控制的必经之路。这两个概念直接决定了数据在多名订阅者之间是如何共享和播放的。我们可以用一个极度接地气的生活场景来建立第一直觉:
- 冷流(默认):就像 “优酷/点播视频”。视频就在服务器放着,不管是张三今天看,还是李四明天看,每个人点开都是从头(第0秒)开始看,大家互不干扰。
- 热流:就像 “电视直播/游戏连麦”。不管有没有人看,直播都在不停地往前推进。张三 12:00 进来看,只能看到 12:00 之后的画面;李四 12:10 进来,前面的 10 分钟他就彻底错过了。
在生产中,我们经常需要把一个“点播流(冷)”转换成“直播流(热)”,或者把直播流“录制缓存(cache)”。这就需要用到 share()、publish() 和 cache() 这三个高级控制配件。
share
share(),多路广播(谁来谁看,不重播)。它的底层是 publish().refCount(1)。它的特点是:当第一个订阅者进来时,水龙头打开,冷流变热流;当所有订阅者都离开(取消订阅)时,流自动关闭。 后来的人只能看到加入那一刻之后的实时数据。
典型场景是聊天室消息广播 / 股市行情推送。股票行情在源源不断地产生。张三和李四看到的是同一时刻的实时价格。李四迟到了 5 秒进来,他不需要、也不应该重新看一遍 5 秒前的历史价格。
1 | public static void main(String[] args) throws InterruptedException { |
1 | 👤 张三订阅了行情... |
publish
publish(),精准控场版变热,手动发车(人齐了再放水)。 share() 虽然能变热,但它有个毛病:第一个人一到就立刻 “放水”,后面来的人注定会漏掉前面的数据。而 publish() 会返回一个特殊的 ConnectableFlux(可连接流)。它就像一辆大巴车,无论张三、李四怎么上车(订阅),大巴车绝对原地不动。直到你作为管理员,手动按下 .connect() 按钮,大巴车才轰然发动,所有上车的人完美同时看到所有数据。
它的典型场景是多系统联动初始化 / 并行拼单抢购。例如启动微服务时,我们需要同时通知“权限模块”和“日志模块”去加载配置项。必须等它们都就位(订阅)了,数据源再统一发射,确保谁也不掉队。
1 | public static void main(String[] args) { |
1 | 📦 模块 1 订阅配置(原地待命)... |
cache
cache(),录像回放版,兼顾实时广播与回放历史。cache() 的威力非常恐怖。它本质上是一个“自带录像机”的热流。 当第一个人订阅时,它不仅开启直播,还会在内存里把直播的内容录下来。当李四迟到进来时,cache() 会做出两步极其惊艳的操作:
- 先把刚才录下来的历史记录,一股脑快进播放(回放)给李四听。
- 历史补齐后,让李四和张三无缝合流,一起看接下来的实时直播。
它的典型场景是行情缓存 / 热门商品详情配置。假设加载一个复杂的 “系统基础字典配置” 需要耗时查库 5 秒。第一个请求进来,触发查库并被 cache() 录制在内存中。后续无论有一万个用户请求进来,都直接从内存录像里秒级读取,不再重复查库。
1 | public static void main(String[] args) throws InterruptedException { |
1 | 👤 张三在第 0 秒订阅... |
Sinks单播/多播/重放
在 Reactor 中,Sinks 的核心作用是显式地在命令式编程(如各种业务 Controller、Service)中作为一个安全的“数据入口”,并将其转化为标准的响应式 Flux 或 Mono。针对不同的并发和分发诉求,Sinks 家族分裂出了三大核心流派:单播(Many().unicast())、多播(Many().multicast()) 和 重放(Many().replay())。
单播
传统管道无法承受“多方抢夺独占资源”。而 Sinks.many().unicast() 创建的流是极致独占的冷流。有且仅允许一个订阅者(Subscriber)进来建立管道连接。如果第二个订阅者敢伸手调用 .subscribe(),它会毫不留情地抛出 IllegalStateException(”UnicastProcessor allows only a single Subscriber”)。如果下游还没来得及订阅,或者处理变慢,进入单播池的数据会暂存在内部的一个内部 FIFO 队列中(默认无界),直到那唯一的订阅者来把它消费完。典型的使用场景包括:
- 独占型后台任务处理:如单线程的文件大批量异步导入、独占的订单流水流水线清洗。
- 严格串行的内部事件分发。
1 | public class UnicastDemo { |
多播
多播解决的是一对多“广播分发”的问题。当一个外部高频事件(如聊天消息、实时股价)产生时,需要同时推送给当前在线的所有人。
- 它是一种热流。它天然允许无数个订阅者同时在线订阅。当一条新数据通过 tryEmitNext 进池时,多播器会像广播喇叭一样,把数据同时复制发送给当前正挂在管道上的所有人。
- 它是 “过时不候” 的。如果数据在上午 10:00 发射,张三在 10:01 订阅进来,那么张三永远看不到 10:00 的那条数据。
典型使用场景包括:
- IM 聊天室/群聊系统:当群主在后台发了一条公告,所有在线的群成员长连接(WebSocket)会同步收到。
- 实时股票行情/大屏看板刷新。
1 | import reactor.core.publisher.Flux; |
1 | 📺 观众 A 看到弹幕: 火箭起飞了! |
重放
解决多播中“迟到的人错过了精彩前情”的痛点。它是带有“记忆功能”的高级多播热流。它不仅具备多播的一对多广播能力,还在内部开辟了一块缓存区(Buffer)。无论订阅者什么时候来,replay() 都会先把缓存区里积攒的 “历史遗留财富” 一股脑全部倾泻给这个新来的人,帮助他完成 “前情提要”,接着再带他并入实时的直播流中。你可以非常精准地控制重放的尺度,比如:重放所有历史数据(all())、只重放最近的 N 个(limit(n))、或者只重放最近一段时间内的(limit(Duration))。
典型使用场景:
- 分布式配置中心推流:当配置项变更时,新启动的微服务订阅进来,必须立刻拿到历史已生效的最新配置。
- IM 聊天历史消息复盘:用户刚进群,不仅能看到实时的聊天,还要自动刷新出最近的 10 条历史聊天记录。
1 | public class ReplayDemo { |
1 | 👤 新成员 [小娟] 刚刚加入了群聊... |
其他问题
Sinks 单播、多播、重放,和 share、publish、cache 的区别是什么?
一句话总结就是:
- Sinks.many().xxx 是【主动泵水门(编程式源头)】:它是一个“热源”,由你手写业务代码主动调用 tryEmitNext 往里塞数据。
- share() / publish() / cache() 是【自来水管分流器(操作符)】:它们本身不是源头,它们必须挂在一个已有的、现成的上游流(通常是冷流)屁股后面,把冷流拦截并转换成热流。
好消息是,Sinks 的三大流派,与这三个操作符在 “对迟到订阅者的态度(缓存策略)” 上,是一一对应的。我们可以把它们两两配对来理解:
- ① 零缓存多播:Sinks…multicast() 和 flux.share()
- 它们的共同点:都是“过时不候”的纯热广播。只把当前最新的数据发给在场的订阅者,迟到的人只能看后面的直播,前面的数据连渣都不剩。
- 它们的区别:
- Sinks…multicast() 是你在外面开个线程,随时 tryEmitNext() 广播。
- flux.share() 是后面挂一个传统的 Flux(比如查数据库)。当第一个人订阅时,上游冷流启动放水;此时如果第二个人也订阅进来,share() 会让第二个人和第一个人共享同一个上游的水流,而不是让数据库重新执行一次查询。
- ② 掌控订阅时机多播:Sinks…multicast() 和 flux.publish()
- 它们的共同点:同样是零缓存的热广播。
- 它们的区别:
- Sinks…multicast() 是第一个人来就立刻放水(自动触发)。
- publish() 更加高冷,它是“不见兔子不撒网”。它返回一个 ConnectableFlux。不管来多少个人订阅,上游都绝对死寂,直到你手动调用了 .connect(),上游才轰隆隆开始放水广播。这能完美保证让所有订阅者“在起跑线上同时听广播”。
- 带记忆的重放:Sinks…replay() 和 flux.cache()
- 它们的共同点:都有极其温情的回放功能。只要有人来(不管迟到多久),都会先把缓存里的历史数据吐给他,再带他看直播。
- 它们的区别:
- Sinks…replay().limit(2) 是你手动塞进去的历史记录。
- flux.cache(2) 则是自动缓存上游冷流冲下来的最后 2 个元素。当后续的新订阅者过来时,不需要重新触发上游冷流(不重新读盘/查库),直接从 cache 操作符的内存里把那 2 个元素吐出来。
注意,在生产环境中,千万不要用 Sinks 去纯粹包装一个已有的冷流,比如写出 Sinks.emit(flux.subscribe()) 这样的套娃代码,这不仅丢失了响应式的背压,还极易引发严重的内存泄漏。
- 遇到 “无中生有,多线程外部塞数据” $\rightarrow$ 用 Sinks 建立根据地。
- 遇到 “已有长流,下游多人争抢防重触发” $\rightarrow$ 用 share()/publish()/cache() 组装分流阀。
响应式编程中的降级和重试
在传统的命令式编程里,异步多线程中的异常处理是个巨大的灾难。如果你在 ExecutorService 的子线程里抛出了未捕获的异常,这个子线程可能会直接暴毙,而主线程甚至毫无察觉。但在响应式世界里,由于 “万物皆为流”,异常(Exception)在底层也被驯化成了一种特殊的“信号”——onError 信号。它顺着自来水管向下游流动,只要我们在管道中间加装 “安全阀”(异常操作符),就能做到不崩掉任何线程池,同时完成优雅降级或自动重试。
在看操作符之前,必须死记一条 Reactive Streams 规范的底层铁律:onError 是一个终止信号。一旦管道中流过了异常信号,这条管道就会彻底宣告报废,后续的数据再也不会流出来了。这也是为什么我们不能用传统的 try-catch 去包住响应式代码,而是必须使用 Reactor 提供的声明式异常操作符。Reactor 提供了三种高频的异常处理手段,就像在水管上装了不同的安全监控,分别是 onErrorReturn、onErrorResume、retry。
onErrorReturn
优雅降级,直接给降级兜底。
1 | Flux.just("user_1", "user_2", "user_error", "user_4") |
1 | # user_4 没有出来。就是因为异常发生后,原管道已经死掉关停了。 |
onErrorResume
偷梁换柱,无缝切到备用管道。如果发生异常,我不只想返回一个死数据,而是想动态启动另一条备用流(比如:查 Redis 挂了,我立刻切换去查本地多级缓存 H2 数据库)。
1 | public static void main(String[] args) throws InterruptedException { |
1 | Redis数据_1 |
retry 或 retryWhen
顽强抵抗,原地无限/有限重试。有些网络抖动是暂时的,盲目降级很可惜。Reactor 提供了惊人的 retry 机制。它的底层原理是,一旦收到 onError 信号,自动重新订阅(resubscribe)最上游的数据源,把管道重新接通再跑一次。
1 | public static void main(String[] args) throws InterruptedException { |
1 | ❌ 尝试发起网络调用... 失败了! |
上述 过程,异常只是在事件循环中作为一个 Throwable 对象在方法入参(onError(t))里传递,没有任何一个物理线程因为这个异常而被迫终止或中断销毁。线程池里的 Worker 线程依然在优哉游哉地处理下一个事件,系统稳如泰山。
庞大的操作符家族
Project Reactor 的核心库 reactor-core 中包含了庞大的操作符家族。如果想完整、权威地查阅所有操作符,唯一的圣经就是其官方的 Javadoc 和参考指南:
- 官方操作符全字典(最推荐的速查表):进入 Reactor 官方参考文档的 Which operator do I need? 交互式向导。
- 官方 API 详细文档(Javadoc):Flux API Javadoc 和 Mono API Javadoc。
将 Reactor(针对 Flux 和 Mono)的所有核心及衍生操作符进行结构化、全景式的分类罗列。

流的创建与源头操作符
这类操作符负责把“各种形态的数据”装进自来水管中,转化为 Flux 或 Mono。
- 基础/静态直接创建
- just(T… data):最简单直接的创建,把内存中的现有对象放入流中。
- empty():创建一个只发出 onComplete 信号的空流。注意和 just(null) 的区别!
- error(Throwable):创建一个只发出 onError 信号的流,用于测试和提前拦截报错。
- never():极其特殊,不发出任何信号(既不给数据,也不结束,也不报错)。
- 传统数据结构/惰性转换
- fromIterable(Iterable) / fromArray(T[]) / fromStream(Stream):把传统的 List、数组、Java 8 Stream 转换为 Flux。
- fromCallable(Callable) / fromRunnable(Runnable):把传统的阻塞/非阻塞同步方法包装成 Mono,被订阅时才真正执行(惰性求值)。
- fromFuture(CompletableFuture):把 Java 原生的异步 Future 无缝接入响应式流。
- defer(Supplier):彻底的懒加载。直到每个订阅者真正调用 subscribe 时,才动态创建全新的流。
- deferContextual:与 defer() 不同,它不捕获创建时的线程上下文,而是在订阅时刻获取当前的 ContextView,确保跨线程(如切换至虚拟线程或线程池)时能正确读取上游写入的上下文 。常用于 WebFlux 链路中透传 TraceID、MDC 日志、安全认证信息等元数据,尤其在涉及 subscribeOn 或虚拟线程调度时防止上下文断裂 。
- 创建动态生成元素的流
- generate(stateSupplier, biFunction):通过 “状态驱动” 的方式,以编程(手写代码)的形式,一个接一个地动态生成元素。
- create(fluxSinkConsumer):它允许你在单次触发中,无拘无束地发射 0 个、1 个、甚至无数个元素(通过 FluxSink),并且完美支持异步、多线程的外部回调。它的终极使命是“连接响应式世界与传统的命令式/事件驱动世界”。如果你手里有一个传统的三方客户端(比如 RocketMQ 的 Consumer 监听器、旧的 WebSocket 监听器、或者一个传统的文件下载进度回调),你想把它们的事件变成 Reactor 的 Flux 响应式流,create 是唯一的正宗选择。
- 时间与定时生成器
- interval(Duration):定时器。每隔一段时间,就从 0 开始发一个递增的 Long 数字(无限流)。
- delayElement(Duration):把流里的元素延迟一段时间再往下投递。
- range(int start, int count):类似于 for 循环,生成一段连续的整数流。
数据的转换与加工
这类操作符用于改变流中元素的形态、类型、或者是流的整体结构。
映射转换(核心中的核心)
- map(Function<T, R>):1对1同步转换。
- mapNotNull(Function):转换时自动过滤掉结果为 null 的元素,会自发把这个 onNext 丢弃掉,不传给下游,防止下游报空指针。。
- flatMap(Function<T, Publisher
>):1对多/1对1异步扁平化转换。将每个元素触发的异步子流融合成一个大流(乱序)。 - flatMapSequential(Function):和 flatMap 一样是异步并发执行,但最终保证输出顺序与输入一致(内部有队列缓存)。
- concatMap(Function):严格按顺序执行的异步转换。处理完前一个元素的子流后,才会开始处理下一个(无并发,保序)。
- switchMap(Function):喜新厌旧的操作符。当新元素进来时,如果前一个元素的子流还没执行完,直接掐断前一个,只听最新的。常用于搜索框联想。
结构重组与归集
- buffer(int maxSize) / buffer(Duration):把散落的元素积攒起来,打包成一个个 List 批量往下送。
- window(int maxSize):和 buffer 类似,但它不打包成 List,而是把元素切块包装成一个个“子 Flux”(流中的流)。
- collectList() / collectMap():把 Flux 中的所有元素全部收集完,最终聚合成一个 Mono<List
> 或 Mono<Map<K, V>>。 - reduce(BiFunction):聚合计算。类似于 Stream 的 reduce,把所有数据累加/累乘/归纳为一个最终结果值。
- scan(BiFunction):和 reduce 一样是聚合,但它会把每一步的中间计算结果都发出来。
转换和加工中的瑞士军刀
- handle(biConsumer, synchronousSink):被称为 “瑞士军刀” 的核心操作符。它完美兼具了 map(转换)和 filter(过滤)的双重超能力,并且能在加工的同时,自主选择是否要提前抛出异常(onError)。
数据的过滤与裁剪
这类操作符负责根据规则,在管道中抛弃某些元素。
条件与边界裁剪
- filter(Predicate):只保留满足条件的元素。
- take(long n) / take(Duration):只拿前 n 个元素(或前一段时间内的数据),拿到后立刻掐断流(自动发送 onComplete)。
- takeLast(int n):只拿流结束前的最后 n 个元素。
- takeUntil(Predicate):一直拿数据,直到某个条件成立的那一刻,停止拿取。
- takeWhile(Predicate):只要条件成立就一直拿,一旦条件不成立立刻停止。
- skip(long n) / skip(Duration):跳过前 n 个元素(或前一段时间内的数据),从后面的开始要。
- defaultIfEmpty(defaultValue):当一个流(Flux 或 Mono)正常结束(发送了 onComplete 信号)时,如果管道里连一个数据(onNext 信号)都没有流过,那么 defaultIfEmpty 就会立刻出手,往管道里塞入一个你指定的默认兜底值送给下游。它只防空流,不防异常。只要流里流过至少一个元素(包括null),这个保底垫子就会完全失效(隐形),绝对不会触发。
- switchIfEmpty(Publisher):入参是一个全新的响应式流(Mono 或 Flux)。适合动态查漏补缺的多级降级方案。比如:“先查本地缓存 $\rightarrow$ 如果为空(switchIfEmpty) $\rightarrow$ 动态发起网络请求去查 Redis $\rightarrow$ 如果还为空(switchIfEmpty) $\rightarrow$ 再去查 MySQL 数据库”。
去重与采样
- distinct():全局去重(底层在内存维护一个 HashSet,注意大数据量时的内存开销)。
- distinctUntilChanged():连续重复去重。只有当当前元素和紧挨着的上一个元素不同时,才允许通过。
- sample(Duration):采样拦截。每隔一段时间,只捞取这段时间内的最后一个元素,其余扔掉。
- ignoreElements():彻底忽略所有数据元素(onNext),只关心流什么时候成功结束(onComplete)或报错(onError),返回一个 Mono
。
流量及背压控制
- limitRate(int prefetch):上游每次最多只会给它灌入 prefetch 个元素,当它把这批数据的 75% 加工完并送给下游后,它就会自动默默地再向上游申请补充剩下的 75% 的货。
- 适用于 “上游是个无脑疯狂发数据的快数据源,而下游处理稍慢” 的情况。如果你不想用 BaseSubscriber 手写背压,直接在中间加一个 limitRate(100),让上游的数据一整批一整批(每批 100 个)地克制下发,极其优雅地保护内存,防止内存溢出(OOM)。
- limitRate(int prefetch, int lowTide):除了指定最大预取量 prefetch(俗称高潮位:High Tide)之外,还允许你打破默认的 75% 规则,亲自指定 “低潮位(Low Tide)”。
- prefetch:第一批向你要多少个。
- lowTide:当我已经消耗了多少个,原有的库存跌落到低潮位时,再去问上游要下一批。
- 如果你希望频繁、小步快跑地补货,把 lowTide 设大一点(比如 limitRate(10, 8),每处理完 2 个就立刻去要 2 个,保持水管内数据持续充沛)。
组合与合并流
当你的数据来自多个不同的源头(比如多个微服务),需要使用这些操作符进行拼装。
并行与交错合并
- zip(Publisher… sources):对齐拼装。把多个流的元素像拉链一样 1对1、2对2 严格对齐,组合成一个 Tuple(元组)下发。如果其中一个流断了,组合立即停止。
- merge(Publisher… sources):无序混流。把多个流汇集到一起,谁的数据先到就先发谁,互不干扰,完全交错。
- combineLatest(Publisher… sources, Function):只要任何一个源头流发出了新数据,就立刻拿这个新数据和其余流中当前“最新的一颗子弹”进行组合输出。
顺序连接
- concat(Publisher… sources):排队连接。先放完第一个流的所有水,再放第二个流的水,严格首尾相接。它是 Flux 类上的一个静态方法。
- concatWith(Publisher… sources):与 concat 的底层核心逻辑是完全一样的,区别仅在于语法层面的表现形式(静态方法 vs 实例链式方法)。
- startWith(T… values):在当前流的数据流出之前,强行在最前面塞入几个固定的先导数据。
线程控制与调度
操控流在哪些物理线程池中横跳。
- publishOn(Scheduler):影响下游执行线程。
- subscribeOn(Scheduler):影响上游源头加载时的线程。
- runOn(Scheduler):专门用于 ParallelFlux(并行流),开启真正的多核多线程并发流水线加工。
错误处理与防御
响应式的 “安全阀门”,负责捕获 onError 信号并进行自救。
- onErrorReturn(T fallbackValue):发生任何/指定异常时,降级返回一个写死的默认值。
- onErrorResume(Function<Throwable, Publisher>):捕获异常,并动态切换到另一个备用流上。
- onErrorMap(Function<Throwable, Throwable>):异常业务翻译。把底层的硬核异常(如 SQLException)转换包装为业务异常(如 UserNotFoundException)再往下抛。
- onErrorComplete():流一旦报错或完成,管道就彻底关闭了。因此,onErrorComplete 虽然让流变成了“正常成功结束”,但报错位置后面的元素依然是出不来的。它只是免去了下游写异常处理回调的麻烦。
- onErrorContinue(biErrorConsumer):当管道在加工某一个元素(比如 map 转换)时不幸报错,onErrorContinue 会拦截错误,把当前这行脏数据扔掉,然后通知管道继续拉取下一个元素加工。注意 onErrorContinue 是一个 “自下而上” 施加影响的全局策略操作符。在编写复杂的响应式链条时,它只能保护在它上游紧邻的、支持它的操作符。如果你把复杂的异步流和不同的调度器(Schedulers)混在一起,使用它必须做充分的单元测试,确保异常的行为如你预期被跳过,否则极易由于线程切换导致跳过策略失效。
- retry(long maxAttempts):一旦报错,无脑原地重新订阅最多 maxAttempts 次。
- retryWhen(Retry):高级重试。可以配置指数退避策略(Backoff),比如第一次隔 100ms 重试,第二次隔 200ms,第三次隔 400ms,防止乱重试把下游系统冲垮。
- timeout:超时设置。
状态监听与生命周期钩子
这类操作符不修改数据流本身,而是像摄像头一样盯着管道的各个节点,通常用来打日志、埋点监控或做清理。
- doOnNext(Consumer):当数据(只包含数据)流过时,偷偷看一下(常用于打印 log.info)。
- doOnEach(Consumer):当数据和信号流过时,偷偷看一下(常用于打印 log.info)。
- doOnError(Consumer):当发生报错时触发,用来打印错误堆栈,但它不会拦截异常,异常依然会向下流。
- doOnSubscribe(Consumer):当下游建立订阅连接的那一刻触发。
- doOnComplete(Runnable):当流正常全部结束时触发。
- doFinally(SignalType):无论流是正常结束、报错崩溃、还是被下游主动 cancel 取消,最终百分之百会执行的代码块(类似于 Java 的 finally,用于清理 IO 资源、关闭连接)
- log():把整个响应式协议(Reactive Streams)在这一层发生的所有交互细节和信号流转,原封不动地打印到你的日志控制台(默认使用 SLF4J 框架)
高级控制与冷热流转换
用于控制流的控制权和多路广播(从冷流走向热流)。
- share():把一个普通的“冷流”(每个订阅者都从头播放)转化为一个“热流”(大家共享同一路实时广播)。
- publish():返回一个 ConnectableFlux。允许你先让成百上千个订阅者在水龙头前排好队,最后调用一次 .connect(),水流同时喷涌而出(实现真正的多路广播控制)。
- cache():缓存流。把流发过的数据缓存在内存里,后续再有新的订阅者进来,不再重新查库,直接把内存里的历史记录播放给它听。
虽然有八大类,但日常开发请务必遵循“用时查表,抓大放小”的原则。绝大多数高并发业务,通过组合 just (创建) $\rightarrow$ filter (裁剪) $\rightarrow$ flatMap/map (转换) $\rightarrow$ zip (组合) $\rightarrow$ onErrorResume (防崩) 这 5 个常规武器,就已经能解决 90% 的生产需求了。
generate 生成流
generate 操作符虽然属于创建这一类,但它与 Flux.just() 或 Flux.fromIterable() 这种 “拿现成数据丢进水管” 的操作符有着本质的区别。generate 是一个 “声明式状态机工厂”。它的作用是通过 “状态驱动” 的方式,以自定义的形式,一个接一个地动态生成元素。
注意,在 generate 的单次循环(或每一轮迭代)中,有且仅能调用一次 synchronousSink.next(data)。也就是说,它每轮只能吐出一个数据。你不能在一次循环里连着调用两次 .next(),也不能不调用(除非你想结束流)。如果你想生成一个无限流(比如斐波那契数列、不断读取网络 Socket、或者动态生成不重复的 UUID),并且这个流的下一个元素严重依赖于上一个元素的状态,那么 generate 就是你的终极武器。
我们先来看一个最基础的案例:从状态 1 开始,每轮乘以 3,直到数字大于 20 为止。
1 | public static void main(String[] args) { |
1 | 当前状态值: 1 -> 乘以3 = 3 |
利用 generate 维护一个 “双变量状态(数组或元组)”,可以极其优雅地吐出一个完美的斐波那契响应式流。
1 | public static void main(String[] args) { |
1 | 0 1 1 2 3 5 8 13 21 34 |
create 生成流
如果说 generate 是一个恪守铁律、只能同步一对一产出的 “严格状态机”,那么 create 就是一个无拘无束、威力无穷的 “全能桥梁适配器”。

create 的典型的使用场景:
- 三方传统监听器(Listener/Callback)的响应式包装:将传统的“触发事件回调”转化为“响应式流数据信号”。
- 多线程并发发射:在外部开启好几个线程并发去拉取数据,然后塞给同一个 FluxSink。
- 自定义复杂的背压缓冲策略:create 允许你在创建流的同时,直接指定当上游发得太快、下游吃不下时的溢出策略(Overflow Strategy),比如 Buffer、Drop、Latest 等。
我们来模拟一个最经典的真实场景:系统里有一个老旧的 UserEventListener(用户登录事件监听器)。每当有用户登录,就会触发 onUserLogin 方法。我们现在用 create 把它无缝包装成响应式的 Flux。
1 | import reactor.core.publisher.Flux; |
1 | 👤 [事件捕获] 用户登录: 张三 |
自定义消费者
虽然普通的 .subscribe(data -> { … }) 传入一个 Consumer 很方便,但它是一个 “全自动” 的水龙头——一旦打开,上游会以排山倒海之势把所有数据一口气推下来(默认请求 Long.MAX_VALUE)。如果你想实现“手控水阀”(比如:下游有空时才分批要数据,或者动态拒绝数据),你就必须通过继承 BaseSubscriber 来自己造一个全功能的“水龙头”。在 Reactor 中,如果你想完全控制消费者的底层行为(比如手动精确控制背压、监听各种生命周期钩子信号),BaseSubscriber 是官方最推荐、最正宗的唯一首选基类。
你可能会想,既然有规范定义的 Subscriber 接口,我自己写个类 implements Subscriber 不行吗?官方不建议这样做。因为直接实现接口,你需要自己去处理复杂的线程安全问题、Subscription 状态维护、以及极其容易发生内存泄漏的并发规则。 而 BaseSubscriber 已经把这些底层的“脏活累活”全部封装好了,并提供了安全的保护机制。
我们来看一个典型的生产案例。上游流水线疯狂流出数据,自定义消费者因为处理能力有限,必须实行严格的 “限流分批” 策略:第一批只要 3 个,每处理完 1 个,才允许上游再送 1 个(动态维持小缓冲区)。
1 | import org.reactivestreams.Subscription; |
1 | 🔌 管道接通! |
Disposable 接口
Disposable 的使用
Disposable 是你调用 .subscribe() 拧开水龙头后,手里拿到的那个 “遥控关阀开关”。当你启动了一个无限流(比如每秒发射一次的定时器,或者一个大文件的读取任务),如果下游在某个时刻不想听了、或者请求已经超时,就可以直接调用 disposable.dispose()。此时下游会向上一路逆流发送一个 cancel(取消)信号通知上游立刻停止生产,从而彻底释放物理线程、网络 Socket 或文件句柄,防止内存泄漏。
1 | public interface Disposable { |
我们来看一个最直观的例子。我们启动一个每隔 500 毫秒就疯狂报时的闹钟流,并在 2 秒后用 Disposable 手动取消它。
1 | public static void main(String[] args) throws InterruptedException { |
1 | ⏰ 闹钟报时: 滴答 0 |
Disposables 复合开关舱
在写响应式项目时,一个用户可能会同时打开好几个数据订阅(比如:接收私聊流、接收群聊流、接收系统通知流)。当这个用户下线(断开 WebSocket)时,我们必须要同时关闭这几百个流。如果一个一个去调用 dispose() 会非常痛苦。Reactor 提供了一个完美的聚合工具类:
1 | import reactor.core.Disposable; |
Context API
Context 的正确姿势
在 Project Reactor 的异步世界里,由于线程经常跨越、横跳(比如从 parallel() 线程池瞬间切到 boundedElastic()),传统的 ThreadLocal 在这里会直接失效——因为数据一旦跨线程,ThreadLocal 就无法自动跟过去了。为了彻底解决 “在整个异步响应式链条中穿透传递共享数据” 的痛点,Reactor 官方祭出了它的终极杀招:Context API(上下文机制)。想要用好 Context,必须先理解以下三条:
- 自下而上(逆流而上)的绑定:这是初学者最容易崩溃的地方。在写代码时,虽然我们是通过 .contextWrite() 来往流里注入数据的,但这个注入动作是在“订阅期(Subscription Time)”自下而上蔓延的。这意味着:一个操作符,只能读取到“在它下游(后面)”注入的 Context 变量。
- 每个订阅者(Subscriber)独享:Context 绑定的是特定的订阅者通道。哪怕同一个 Flux 管道被并发调用了 1 万次,这 1 万个请求的 Context 也是完全隔离、互不干扰的。
- 不可变性(Immutable):Context 的底层设计和 Java 的 String 一样,是完全不可变的。每次你调用 context.put(k, v),底层都会克隆并返回一个新的 Context 对象。这在多线程并发时天然保证了线程安全。
错误示范:把注入写在了读取的“上游”
1 | Mono.just("Hello") |
正确示范:把注入写在读取的“下游”(后面)
1 | Mono.just("Hello") |
简单案例:
1 | public static void main(String[] args) { |
1 | 用户的 HTTP 请求进入 Controller... |
使用新版 API
由于 Reactor 的演进,旧版本的 mono.subscriberContext() 已经被彻底废弃。现代开发请使用以下两组正宗 API:
- ① 往流里“写”数据(操作符)
- .contextWrite(Context):直接覆盖或注入一个现成的 Context。
- .contextWrite(ctx -> ctx.put(“K”, “V”)):通过 Lambda 表达式,在原有的 Context 基础上追加键值对。
- ② 从流里“读”数据(源头/中游)
- Mono.deferContextual(ctx -> …):高频首选。在流的任意位置,如果需要用到 Context,直接用它包裹住你的业务逻辑。
- Flux.deferContextual(ctx -> …):同上,适用于多元素流。
需要注意的地方:
- 跟 Schedulers 混用时放最后:在复杂的微服务架构中,为了防止线程池横跳导致 Context 信号在中途丢失,最保险的做法是将 .contextWrite() 尽量写在靠近 .subscribe()(流的最末端)的位置。
- 它不是万能的缓存(Cache):不要把业务大对象(如整个 User 实体、大 List)塞进 Context 里,它在底层每次 put 都会克隆。它只适合存放轻量级的元数据:如 TraceId、TenantId(租户ID)、UserToken 字符串、安全认证的 Roles 列表等。
用 Reactor API 写 Server
一个最简单的示例
如果我们完全脱离 Spring Boot、脱离 Spring WebFlux,只用最原生的 Reactor API 加上底层的驱动网络库来搭建一个高性能的 Web 服务器,我们需要使用 Project Reactor 官方专为网络通信打造的组件:Reactor Netty。它是一个基于 Netty 的响应式网络库。它直接将 Netty 的通道(Channel)和事件循环(EventLoop)抽象成了 Reactor 的 Mono 和 Flux。
1 | import reactor.core.publisher.Mono; |
直接右键运行上方的 main 方法,打开浏览器或者使用 curl 访问:
1 | curl http://localhost:8088/hello |
测试异步非阻塞 POST 响应式流:
1 | curl -X POST -d "koohub koo code, cool webflux" http://localhost:8088/echo |
模拟 ChatGPT 流式响应
1 | package com.zdemo.scloud; |
在项目 resources 目录,创建一个名为 index.html 的文件。传统的 EventSource 只支持 GET 请求。为了配合 AI 业务中常见的 POST 提问,前端使用现代浏览器的 fetch API 和 ReadableStream(可读流) 来动态拦截并渲染后端的响应式流。
1 |
|
联调和效果展示:
- 启动后端:运行 PureReactorServer 的 main 方法。
- 打开前端:在浏览器直接输入并打开地址:http://localhost:8088/ 或者是双击本地的 index.html。
- 点击“发射提问”:你会看到下方的黑色控制台区域,完全没有长时间的空白等待,而是像 ChatGPT 官方一样,文字以每秒 20 个字的高频速度,如丝般顺滑地“吐”在屏幕上!