Java 响应式编程 - Reactor 基础库

学习 Reactor 绝对不能盲从,因为网上有大量旧版本或错误的旧观念。需要我们死磕官方一手资料:


学习路线

两条管道 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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
<dependencyManagement>
<dependencies>
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-bom</artifactId>
<version>${reactor-bom.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>

<dependencies>
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-core</artifactId>
</dependency>
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>

如果你当前的系统是标准的 Spring Boot 项目(如使用了 WebFlux),那么你完全不需要手动去指定 reactor-bom 的版本。因为 Spring Boot 已经在它自己的官方 spring-boot-dependencies 中做好了最佳的版本依赖仲裁。你只需要像下面这样直接引入 WebFlux 或者基础的 starter,它会自动给你带上最兼容、最稳定的对应 Reactor 版本:

1
2
3
4
5
6
7
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
<version>3.3.0</version>
</dependency>
</dependencies>


创建管道

在 Reactor 中,代码的第一步永远是“把数据塞进管道”,变成 Flux 或 Mono。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.util.Arrays;
import java.util.List;

public class QuickStart {
public static void main(String[] args) {
// 1. 创建 Mono(0 或 1 个元素):就像单发子弹
Mono<String> monoSource = Mono.just("唯一的数据");

// 2. 创建 Flux(0 到 N 个元素):就像机关枪弹链
Flux<String> fluxSource1 = Flux.just("数据A", "数据B", "数据C");

// 3. 从现有的传统 List 集合创建 Flux 管道
List<String> list = Arrays.asList("张三", "李四", "王五");
Flux<String> fluxSource2 = Flux.fromIterable(list);
}
}

此时的重点:如果你直接运行上面的代码,控制台什么都不会输出。 因为你现在只是画好了“自来水管的设计图”。记住那句至理名言:不订阅,什么都不会发生


订阅通水

要想让水管里的水流出来,必须在管道的最末端安装一个水龙头,也就是调用 subscribe()。

1
2
3
4
5
6
7
8
9
10
11
Flux<String> flux = Flux.just("Java", "Python", "Go");

// 最基础的订阅(只接水)
flux.subscribe(data -> System.out.println("收到: " + data));

// 全功能订阅(接水、管报错、管结束)
flux.subscribe(
data -> System.out.println("【正常数据】: " + data), // 1. 对应 onNext
error -> System.err.println("【发生暴雷】: " + error), // 2. 对应 onError
() -> System.out.println("【全部流完,完美闭环!】") // 3. 对应 onComplete
);


初试管道配件

现在水能流过去了,我们要在水管中间加装 “净水器”。这就是操作符。

1
2
3
4
5
6
7
8
9
Flux<Integer> scoreFlux = Flux.just(45, 78, 92, 59, 100);

scoreFlux
// 1. 拦截过滤:只要大于等于 60 分的(及格线)
.filter(score -> score >= 60)
// 2. 转换加工:把数字类型的分数,包装成一句话
.map(score -> "合格分数: " + score)
// 3. 安装水龙头放水
.subscribe(System.out::println);


异步编排的灵魂 flatMap

在 Project Reactor 的几百个操作符中,flatMap 是最具分水岭意义的核心操作符,也是初学者最容易卡壳的地方。在命令式编程中,如果我们要对一组数据进行 “再去查一次数据库/调一次微服务” 的操作,通常会用 for 循环里套一个 RPC 调用。但在响应式世界里,这种做法是致命的。要想彻底精通 flatMap,我们先从它的兄弟 map 开始对比,用自来水管的哲学来一击看穿。用一句话概括它们的不同:

  • map:同步转换。进去一个元素,出来一个普通元素(一对一)。
  • flatMap:异步扁平化转换。进去一个元素,出来一个全新的子管道(Flux/Mono),并把所有子管道融合成一条主流。

如果我们错误地使用 map,会发生什么?

1
2
3
4
5
6
7
// 假设这是一个异步查库的方法,返回一个 Mono 管道
public Mono<String> getUserNameFromDb(Integer userId) {
return Mono.just("用户-" + userId).delayElement(Duration.ofMillis(100));
}
// 错误的尝试:使用 map
Flux<Integer> userIdFlux = Flux.just(1, 2, 3);
Flux<Mono<String>> resultWithMap = userIdFlux.map(id -> getUserNameFromDb(id));

因为 getUserNameFromDb 本身返回的是一个管道(Mono),map 只是机械地把元素包进去。结果导致你得到的是一个 “套娃管道” (Flux)。你装水龙头(subscribe)的时候,流出来的不是水(数据),而是一个个套着的水管!你根本拿不到里面的数据。这时候就需要使用到 flatMap!

1
2
3
Flux<Integer> userIdFlux = Flux.just(1, 2, 3);
Flux<String> resultWithFlatMap = userIdFlux.flatMap(id -> getUserNameFromDb(id));
resultWithFlatMap.subscribe(System.out::println);

底层流转轨迹:

  • 分流:当元素 1 流过 flatMap 时,flatMap 调用 getUserNameFromDb(1),产生了一个独立的、异步的子管道 Mono。
  • 并发执行:元素 2、3 过来,同样各自产生了自己的异步子管道。这些子管道在底层同时(并行)开始执行异步查库。
  • 扁平化融合:flatMap 内部充当了一个 “总集水器”。哪个子管道的数据先从数据库返回,它就把这滴 “水” 捞出来,顺着主管道送给最终的订阅者。

下面看一个简单的例子:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
public static void main(String[] args) throws InterruptedException {
long startTime = System.currentTimeMillis();
Flux<String> orderIdFlux = Flux.just("Order_A", "Order_B", "Order_C");

orderIdFlux.flatMap(FlatMapAdvanced::getOrderDetails)
.subscribe(
data -> System.out.println("🔥 ["+ Thread.currentThread().getName() +"]消费者收到 -> " + data + " | 耗时: " + (System.currentTimeMillis() - startTime) + "ms"),
err -> {},
() -> System.out.println("🎉 ["+ Thread.currentThread().getName() +"]所有异步任务处理完毕!总耗时: " + (System.currentTimeMillis() - startTime) + "ms")
);

Thread.sleep(3000);
}

/**
* 模拟一个极其耗时的异步订单详情查询操作(每个订单耗时 1 秒)
*/
private static Mono<String> getOrderDetails(String orderId) {
return Mono.just("[create in "+ Thread.currentThread().getName() +"]订单详情: " + orderId)
.delayElement(Duration.ofSeconds(1)); // 模拟异步网络延迟 1 秒
}
1
2
3
4
🔥 [parallel-1]消费者收到 -> [create in main]订单详情: Order_B | 耗时: 1085ms
🔥 [parallel-1]消费者收到 -> [create in main]订单详情: Order_A | 耗时: 1086ms
🔥 [parallel-1]消费者收到 -> [create in main]订单详情: Order_C | 耗时: 1086ms
🎉 [parallel-1]所有异步任务处理完毕!总耗时: 1087ms
  • 总耗时只有 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
2
3
4
🔥 [parallel-1]消费者收到 -> [create in main]订单详情: Order_A | 耗时: 1492ms
🔥 [parallel-2]消费者收到 -> [create in parallel-1]订单详情: Order_B | 耗时: 2500ms
🔥 [parallel-3]消费者收到 -> [create in parallel-2]订单详情: Order_C | 耗时: 3503ms
🎉 [parallel-3]所有异步任务处理完毕!总耗时: 3504ms


自定义hanle操作

在日常开发中,我们经常遇到这样的骚操作:拿到一个数据 $\rightarrow$ 先判断它合不合规(filter) $\rightarrow$ 如果合规,把它转成另一种对象(map)。如果用常规写法,你需要连着写 .filter(…).map(…)。但如果使用 handle,它会提供给你一个全能的 SynchronousSink(同步发射器),让你在一个方法块里自发决定这个元素是 “保留并转换”、“直接扔掉”还是“拉响警报(报错)”。注意,虽然它很自由,但它和 generate 一样,在单次处理某一个元素时,最多只能调用一次 sink.next(data)(可以不调用,代表过滤掉;调用一次,代表加工传给下游)。

典型的使用场景,如全能数据清洗仓(过滤 + 转换 + 报错)。假设你正在负责某后台用户评论数据清洗。你需要对一串文章 ID 或者是内容流进行处理:

  • 如果内容包含“违规恶意政治”,立刻拉响警报(抛出异常,触发降级)。
  • 如果内容包含“广告/水贴”,直接无视它(过滤掉)。
  • 如果内容合法,自动给它追加“官方审核”的后缀(转换)。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
public class HandleDemo {
public static void main(String[] args) {

// 模拟一串各种各样的帖子文本流
Flux<String> postFlux = Flux.just("正常内容A", "恶意政治垃圾", "水贴111", "正常内容B");

Flux<String> processedFlux = postFlux.handle((post, sink) -> {
// 场景一:发现极度危险的内容,直接拉响警报,中断管道
if (post.contains("政治")) {
sink.error(new RuntimeException("🚨 监控到恶意违规政治内容,立即拦截报警!"));
return; // 拦截,不往下走
}

// 场景二:发现水贴广告,直接忽略(不调用 sink.next 相当于被 filter 过滤掉了)
if (post.contains("水贴")) {
System.out.println(" [handle内部] 发现水贴: [" + post + "],已自动做丢弃处理。");
return;
}

// 场景三:合法内容,进行加工转换并下发给下游(相当于 map)
sink.next("📝 " + post + " [已通过系统合规审核]");
});

// 订阅放水
processedFlux.subscribe(
System.out::println,
err -> System.err.println(err.getMessage())
);
}
}
1
2
📝 正常内容A [已通过系统合规审核]
🚨 监控到恶意违规政治内容,立即拦截报警!


transform [Deferred]

这两者的本质区别是它们对外部变量的“捕获时机”不同,而这个时机差导致了“是否能获取到外部变量最新值”的本质区别。

  • transform(Function) —— 组装期捕获(只看应用启动那一刻)。当你的代码执行到 transform 这一行时(通常是 Spring 容器启动、对象初始化、或者管道声明时),Reactor 就会立刻执行里面的 Function。如果 Function 内部去读取了一个外部变量,它在组装期就把这个变量的值给固定下来了。哪怕这个外部变量是个 AtomicInteger 或者普通的成员变量,之后它的值在运行时变了,后续的所有订阅者,看到的依然是组装期捕获的旧值。
  • transformDeferred(Function) —— 订阅期捕获(千人千面,看拧开水龙头那一刻)。代码执行到 transformDeferred 时,它在组装期什么都不做(延迟执行)。只有当某个请求过来说 flux.subscribe()(拧开水龙头)时,它才会当场现去执行里面的 Function。这时 Function 才会去读取外部变量。因此它每次都能抓取到外部变量在 “当下这一秒” 的最新状态。

我们来看下面这个例子:

1
2
3
4
5
6
7
8
9
10
11
12
13
public static void main(String[] args) {
AtomicInteger i = new AtomicInteger(0);
Flux<String> demoFlux = Flux.just("a", "b", "c")
.transform(flux -> {
if (i.incrementAndGet() == 1) {
return flux.map(String::toUpperCase);
} else {
return flux;
}
});
demoFlux.subscribe(v -> System.out.println("订阅者1;v=" + v));
demoFlux.subscribe(v -> System.out.println("订阅者2;v=" + v));
}
1
2
3
4
5
6
订阅者1;v=A
订阅者1;v=B
订阅者1;v=C
订阅者2;v=A
订阅者2;v=B
订阅者2;v=C

将上面的 transform 改成 transformDefer,执行结果就会变成:

1
2
3
4
5
6
订阅者1;v=A
订阅者1;v=B
订阅者1;v=C
订阅者2;v=a
订阅者2;v=b
订阅者2;v=c

其实,这两兄弟的诞生,不是为了对流里的数据做具体的加减乘除,而是为了解决响应式开发中的一个终极痛点:代码如何优雅地实现复用(响应式切面重构)。一句话切中它们的本质区别:

  • transform 属于“静态组装”:在代码编译/编译后的链条组装期,逻辑就已经写死了,且全局只执行一次。
  • transformDeferred 属于“动态加载”:在每一个订阅者真正调用 subscribe() 的订阅期(Subscription Time)才会触发执行,每个订阅者都会单独动态执行一次。

再来看一个使用 transform 重构的示例:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
public class TransformDemo {
public static void main(String[] args) {
// 1. 定义一个公共的高频复用“零件包”(把一个旧 Flux 包装并返回一个新 Flux)
Function<Flux<String>, Flux<String>> addSignatureAndFilter =
flux -> flux.filter(val -> val != null && !val.isEmpty())
.map(val -> "[安全认证] " + val);

// 2. 在不同的业务管道里,直接用 transform 一键组装,免去重写一大堆 filter 和 map
Flux.just("数据A", "数据B")
.transform(addSignatureAndFilter) // 🌟 静态组装进管道
.subscribe(System.out::println);

Flux.just("数据C")
.transform(addSignatureAndFilter)
.subscribe(System.out::println);
}
}

总结起来,只要你的重构逻辑需要动态适配外部配置变更、动态读取本地缓存、或者动态追踪线程上下文(如 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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
public static void main(String[] args) throws InterruptedException {
Flux.just("user_101", "user_102") // 假设最初由 Netty 线程(Thread-A)触发

// 【当前在 Thread-A】
.map(id -> {
System.out.println("✏️ ["+ Thread.currentThread().getName() +"] 提取ID: " + id);
return id;
})

// 🚪 传送门 1:切换到 Parallel 线程池
.publishOn(Schedulers.parallel())

// 【当前在 parallel-x 线程】
.map(id -> {
System.out.println("🔐 ["+ Thread.currentThread().getName() +"] 复杂的加密计算...");
return "ENCRYPTED_" + id;
})

// 🚪 传送门 2:切换到 BoundedElastic 线程池(专门容纳阻塞操作)
.publishOn(Schedulers.boundedElastic())

// 【当前在 boundedElastic-x 线程】
.map(id -> {
try {
System.out.println("🛢️ ["+ Thread.currentThread().getName() +"] 执行传统阻塞 JDBC 查库...");
Thread.sleep(200); // 模拟传统的阻塞延时
} catch (InterruptedException e) {}
return id + "_DB_RESULT";
})

// 安装水龙头,通水!
.subscribe(data -> System.out.println("🎉 ["+ Thread.currentThread().getName() +"] 最终消费: " + data));

// 等待异步流执行完
Thread.sleep(2000);
}
1
2
3
4
5
6
7
8
✏️ [main] 提取ID: user_101
✏️ [main] 提取ID: user_102
🔐 [parallel-1] 复杂的加密计算...
🔐 [parallel-1] 复杂的加密计算...
🛢️ [boundedElastic-1] 执行传统阻塞 JDBC 查库...
🎉 [boundedElastic-1] 最终消费: ENCRYPTED_user_101_DB_RESULT
🛢️ [boundedElastic-1] 执行传统阻塞 JDBC 查库...
🎉 [boundedElastic-1] 最终消费: ENCRYPTED_user_102_DB_RESULT
  • 初始流由 main 线程触发(在 WebFlux 中则是 Netty 线程)。
  • 当数据撞击到第一个 publishOn(Schedulers.parallel()) 时,底层的无锁环形队列发挥作用,数据被递交给 parallel-1 线程,随后的加密 map 立即在计算池中运行。
  • 当再次撞击到 publishOn(Schedulers.boundedElastic()) 时,空间再次切换。传统阻塞查库被安放在了专门的弹性阻塞池中,完美保护了主计算池和网络池不被卡死。

为什么说 subscribeOn 是“逆流起效”的?我们来看下面这个案例:

1
2
3
4
5
6
7
8
9
10
11
public static void main(String[] args) throws InterruptedException {
Flux.just("user_101", "user_102") // 假设最初由 Netty 线程(Thread-A)触发
.subscribeOn(Schedulers.boundedElastic()) // 👈 写在后面
.map(data -> {
System.out.println("["+ Thread.currentThread().getName() +"] 加工数据");
return data;
}).subscribe();

// 等待异步流执行完
Thread.sleep(2000);
}
1
2
[boundedElastic-1] 加工数据
[boundedElastic-1] 加工数据
  • 订阅期(从下往上):调用 subscribe() 时,信号像多米诺骨牌一样从下往上逆流回溯。
  • 触发 subscribeOn:往上回溯遇到 subscribeOn 时,它强行接管:“别往上找了,从现在开始,到最顶端数据源加载的这个逆流动作,由我指定的 boundedElastic-1 线程来干!”
  • 运行时(从上往下)**:于是,最顶层的数据源在 boundedElastic-1 线程中被加载,并顺流而下,一路带偏了后面的 map。

这样做的好处是,即使这个数据库查询要耗时 5 秒,由于使用了 subscribeOn(Schedulers.boundedElastic()),整个阻塞过程也只会占用 boundedElastic 池子里的线程。而 WebFlux 核心的那几个极其宝贵的、负责接网卡 HTTP 请求的 Netty Worker 线程依然一身轻松,可以继续疯狂接收其他用户的请求。这就是 Reactor 的“幻影移形”魔法,利用空间换时间,用精细的线程池隔离,达成了整体系统的高可用。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
public static void main(String[] args) throws InterruptedException {
Flux.just("user_101", "user_102")
.subscribeOn(Schedulers.boundedElastic())
.map(data -> {
System.out.println("["+ Thread.currentThread().getName() +"] 加工数据1");
return data;
})
.publishOn(Schedulers.parallel())
.map(data -> {
System.out.println("["+ Thread.currentThread().getName() +"] 加工数据2");
return data;
})
.subscribe();
Thread.sleep(2000);
}
1
2
3
4
[boundedElastic-1] 加工数据1
[boundedElastic-1] 加工数据1
[parallel-1] 加工数据2
[parallel-1] 加工数据2


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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
import reactor.core.publisher.Flux;
import reactor.core.scheduler.Schedulers;

public class ParallelRunOnDemo {
public static void main(String[] args) throws InterruptedException {
long startTime = System.currentTimeMillis();

// 1. 准备 8 个需要加密的密码
Flux.just("pwd1", "pwd2", "pwd3", "pwd4", "pwd5", "pwd6", "pwd7", "pwd8")

// 2. 核心大招:开启并行流,默认切分成等同于 CPU 核心数的 “轨道”
.parallel()

// 3. 为每个轨道分配独立的物理线程池(计算密集池)
.runOn(Schedulers.parallel())

// 4. 在各自的轨道上,并发执行重度 CPU 计算
.map(pwd -> heavyBCryptHash(pwd))

// 5. 聚合:把并行的各个轨道重新融合成一个主流 Flux
.sequential()

// 6. 订阅放水
.subscribe(
hash -> System.out.println("🔑 " + hash),
err -> {},
() -> System.out.println("\n🎉 所有密码处理完毕!总耗时: " + (System.currentTimeMillis() - startTime) + "ms")
);

// 异步等待
Thread.sleep(2000);
}

/**
* 模拟极其消耗 CPU 的强哈希加密算法(耗时 100ms)
*/
private static String heavyBCryptHash(String password) {
try {
Thread.sleep(100);
} catch (InterruptedException e) {}
return "HASHED_[" + password + "]_BY_" + Thread.currentThread().getName();
}
}
1
2
3
4
5
6
7
8
9
10
🔑 HASHED_[pwd1]_BY_parallel-1
🔑 HASHED_[pwd2]_BY_parallel-2
🔑 HASHED_[pwd3]_BY_parallel-3
🔑 HASHED_[pwd4]_BY_parallel-4
🔑 HASHED_[pwd5]_BY_parallel-1
🔑 HASHED_[pwd6]_BY_parallel-2
🔑 HASHED_[pwd7]_BY_parallel-3
🔑 HASHED_[pwd8]_BY_parallel-4

🎉 所有密码处理完毕!总耗时: 218ms
  • 耗时被极度压缩:总共 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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
public static void main(String[] args) throws InterruptedException {
// 创建一个冷流:每秒递增一个数字
Flux<Long> liveTicker = Flux.interval(Duration.ofSeconds(1))
.doOnSubscribe(s -> System.out.println("▶️ 原始水源被激活!"))
.share(); // 🌟 核心:一键切成热流(广播模式)

// 订阅者1:第 0 秒加入
System.out.println("👤 张三订阅了行情...");
liveTicker.subscribe(data -> System.out.println(" [张三] 看到实时价格: " + data));

Thread.sleep(2500); // 过去 2.5 秒

// 订阅者2:第 2.5 秒迟到加入
System.out.println("👤 李四迟到加入...");
liveTicker.subscribe(data -> System.out.println(" [李四] 看到实时价格: " + data));

Thread.sleep(30000);
}
1
2
3
4
5
6
7
8
9
👤 张三订阅了行情...
▶️ 原始水源被激活!
[张三] 看到实时价格: 0
[张三] 看到实时价格: 1
👤 李四迟到加入...
[张三] 看到实时价格: 2
[李四] 看到实时价格: 2 <-- 💡 李四没有从 0 开始看,而是直接和张三同步看 2!
[张三] 看到实时价格: 3
[李四] 看到实时价格: 3


publish

publish(),精准控场版变热,手动发车(人齐了再放水)。 share() 虽然能变热,但它有个毛病:第一个人一到就立刻 “放水”,后面来的人注定会漏掉前面的数据。而 publish() 会返回一个特殊的 ConnectableFlux(可连接流)。它就像一辆大巴车,无论张三、李四怎么上车(订阅),大巴车绝对原地不动。直到你作为管理员,手动按下 .connect() 按钮,大巴车才轰然发动,所有上车的人完美同时看到所有数据。

它的典型场景是多系统联动初始化 / 并行拼单抢购。例如启动微服务时,我们需要同时通知“权限模块”和“日志模块”去加载配置项。必须等它们都就位(订阅)了,数据源再统一发射,确保谁也不掉队。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
public static void main(String[] args) {
Flux<String> configSource = Flux.just("配置项A", "配置项B", "配置项C");

// 🌟 核心:转化成 ConnectableFlux,此时处于“待命”状态
ConnectableFlux<String> published = configSource.publish();

System.out.println("📦 模块 1 订阅配置(原地待命)...");
published.subscribe(data -> System.out.println(" [模块1] 加载: " + data));

System.out.println("📦 模块 2 订阅配置(原地待命)...");
published.subscribe(data -> System.out.println(" [模块2] 加载: " + data));

System.out.println("🚀 所有人已准备就绪,管理员按下发车键 (.connect())!\n");
published.connect(); // 🌟 核心动作:千军万马同时奔流!
}
1
2
3
4
5
6
7
8
9
10
📦 模块 1 订阅配置(原地待命)...
📦 模块 2 订阅配置(原地待命)...
🚀 所有人已准备就绪,管理员按下发车键 (.connect())!

[模块1] 加载: 配置项A
[模块2] 加载: 配置项A
[模块1] 加载: 配置项B
[模块2] 加载: 配置项B
[模块1] 加载: 配置项C
[模块2] 加载: 配置项C


cache

cache(),录像回放版,兼顾实时广播与回放历史。cache() 的威力非常恐怖。它本质上是一个“自带录像机”的热流。 当第一个人订阅时,它不仅开启直播,还会在内存里把直播的内容录下来。当李四迟到进来时,cache() 会做出两步极其惊艳的操作:

  • 先把刚才录下来的历史记录,一股脑快进播放(回放)给李四听。
  • 历史补齐后,让李四和张三无缝合流,一起看接下来的实时直播。

它的典型场景是行情缓存 / 热门商品详情配置。假设加载一个复杂的 “系统基础字典配置” 需要耗时查库 5 秒。第一个请求进来,触发查库并被 cache() 录制在内存中。后续无论有一万个用户请求进来,都直接从内存录像里秒级读取,不再重复查库。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
public static void main(String[] args) throws InterruptedException {
// 创建一个冷流,并开启缓存录制(默认缓存无限个,也可以指定 cache(2) 缓存最近2个)
Flux<Long> cachedFlow = Flux.interval(Duration.ofSeconds(1))
.doOnSubscribe(s -> System.out.println("🛢️ 查库/耗时源头被触发!"))
.cache(); // 也可以执指定 history 个数

System.out.println("👤 张三在第 0 秒订阅...");
cachedFlow.subscribe(data -> System.out.println(" [张三]: " + data));

Thread.sleep(2500); // 过去 2.5 秒,流已经发出了 0, 1

System.out.println("\n👤 李四在第 2.5 秒迟到加入...");
// 💡 奇迹发生:李四进来的瞬间,会立刻收到历史录像 0, 1
cachedFlow.subscribe(data -> System.out.println(" [李四]: " + data));

Thread.sleep(20000);
}
1
2
3
4
5
6
7
8
9
10
11
12
👤 张三在第 0 秒订阅...
🛢️ 查库/耗时源头被触发!
[张三]: 0
[张三]: 1

👤 李四在第 2.5 秒迟到加入...
[李四]: 0 <-- 🎥 历史回放
[李四]: 1 <-- 🎥 历史回放
[张三]: 2 <-- 🛠️ 重新合流,一起看最新的实时直播
[李四]: 2
[张三]: 3
[李四]: 3


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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
public class UnicastDemo {
public static void main(String[] args) {
// 1. 创建一个单播 Sink
Sinks.Many<String> logSink = Sinks.many().unicast().onBackpressureBuffer();

// 2. 模拟系统在各个地方、由不同线程异步泵入日志事件
logSink.tryEmitNext("Log-1: 用户登录成功");
logSink.tryEmitNext("Log-2: 支付发起");

// 3. 唯一的消费者:后台主日志写入器建立连接
Flux<String> logFlux = logSink.asFlux();
logFlux.subscribe(log -> System.out.println("[写入磁盘] " + log));

// ❌ 危险操作:如果有第二个消费者想插队订阅,直接暴毙
// logFlux.subscribe(log -> System.out.println("抢夺流量"));
}
}


多播

多播解决的是一对多“广播分发”的问题。当一个外部高频事件(如聊天消息、实时股价)产生时,需要同时推送给当前在线的所有人。

  • 它是一种热流。它天然允许无数个订阅者同时在线订阅。当一条新数据通过 tryEmitNext 进池时,多播器会像广播喇叭一样,把数据同时复制发送给当前正挂在管道上的所有人。
  • 它是 “过时不候” 的。如果数据在上午 10:00 发射,张三在 10:01 订阅进来,那么张三永远看不到 10:00 的那条数据。

典型使用场景包括:

  • IM 聊天室/群聊系统:当群主在后台发了一条公告,所有在线的群成员长连接(WebSocket)会同步收到。
  • 实时股票行情/大屏看板刷新。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
import reactor.core.publisher.Flux;
import reactor.core.publisher.Sinks;

public class MulticastDemo {
public static void main(String[] args) throws InterruptedException {
// 1. 创建一个多播 Sink(配置:如果暂时没人订阅,数据自动丢弃,保持热流本色)
Sinks.Many<String> liveSink = Sinks.many().multicast().onBackpressureBuffer();
Flux<String> liveFlux = liveSink.asFlux();

// 2. 观众 A 率先入场,开始看直播
liveFlux.subscribe(data -> System.out.println("📺 观众 A 看到弹幕: " + data));

// 3. 现场发射前两条弹幕
liveSink.tryEmitNext("火箭起飞了!");
liveSink.tryEmitNext("666666");

Thread.sleep(100);

// 4. 观众 B 迟到了,刚刚连接进来
System.out.println("\n👤 观众 B 慢悠悠地进入了直播间...\n");
liveFlux.subscribe(data -> System.out.println("📱 观众 B 看到弹幕: " + data));

// 5. 现场再次发射一条新弹幕
liveSink.tryEmitNext("主播真帅!");
}
}
1
2
3
4
5
6
7
📺 观众 A 看到弹幕: 火箭起飞了!
📺 观众 A 看到弹幕: 666666

👤 观众 B 慢悠悠地进入了直播间...

📺 观众 A 看到弹幕: 主播真帅!
📱 观众 B 看到弹幕: 主播真帅! <-- 迟到的观众 B 只能看到进来后的新弹幕,前面的错过了


重放

解决多播中“迟到的人错过了精彩前情”的痛点。它是带有“记忆功能”的高级多播热流。它不仅具备多播的一对多广播能力,还在内部开辟了一块缓存区(Buffer)。无论订阅者什么时候来,replay() 都会先把缓存区里积攒的 “历史遗留财富” 一股脑全部倾泻给这个新来的人,帮助他完成 “前情提要”,接着再带他并入实时的直播流中。你可以非常精准地控制重放的尺度,比如:重放所有历史数据(all())、只重放最近的 N 个(limit(n))、或者只重放最近一段时间内的(limit(Duration))。

典型使用场景:

  • 分布式配置中心推流:当配置项变更时,新启动的微服务订阅进来,必须立刻拿到历史已生效的最新配置。
  • IM 聊天历史消息复盘:用户刚进群,不仅能看到实时的聊天,还要自动刷新出最近的 10 条历史聊天记录。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
public class ReplayDemo {
public static void main(String[] args) {
// 1. 创建一个只重放最近 2 个元素的重放 Sink
Sinks.Many<String> chatSink = Sinks.many().replay().limit(2);
Flux<String> chatFlux = chatSink.asFlux();

// 2. 房间里零星发生了 3 次对话
chatSink.tryEmitNext("消息1:大家早上好!");
chatSink.tryEmitNext("消息2:今天 Owlias 博客更新 v1.3 了。");
chatSink.tryEmitNext("消息3:增加了很多好玩的功能。");

// 3. 这时候,一个新成员刚加群订阅进来
System.out.println("👤 新成员 [小娟] 刚刚加入了群聊...");
chatFlux.subscribe(msg -> System.out.println(" [历史/实时消息] " + msg));

// 4. 群里发了一条新消息
chatSink.tryEmitNext("消息4:欢迎新成员!");
}
}
1
2
3
4
👤 新成员 [小娟] 刚刚加入了群聊...
[历史/实时消息] 消息2:今天 Owlias 博客更新 v1.3 了。 <-- 🌟 自动向前追溯重放历史
[历史/实时消息] 消息3:增加了很多好玩的功能。 <-- 🌟 自动向前追溯重放历史
[历史/实时消息] 消息4:欢迎新成员! <-- ⚡ 接轨之后的实时多播


其他问题

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
2
3
4
5
6
7
8
9
10
Flux.just("user_1", "user_2", "user_error", "user_4")
.map(user -> {
if ("user_error".equals(user)) {
throw new RuntimeException("数据库连接超时!");
}
return user + "_已认证";
})
// 安全阀 1:一旦前面发生任何异常,立刻切断原流,并向下游发送一个“匿名用户”作为兜底值
.onErrorReturn("匿名用户(系统降级)")
.subscribe(System.out::println);
1
2
3
4
# user_4 没有出来。就是因为异常发生后,原管道已经死掉关停了。
user_1_已认证
user_2_已认证
匿名用户(系统降级)


onErrorResume

偷梁换柱,无缝切到备用管道。如果发生异常,我不只想返回一个死数据,而是想动态启动另一条备用流(比如:查 Redis 挂了,我立刻切换去查本地多级缓存 H2 数据库)。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
public static void main(String[] args) throws InterruptedException {
Flux<String> primaryRedisFlow = Flux.just("Redis数据_1", "Redis数据_2")
.concatWith(Flux.error(new RuntimeException("Redis 集群崩了!"))) // 模拟突发崩溃
.concatWith(Flux.just("Redis数据_3"));
Flux<String> backupLocalCacheFlow = Flux.just("本地缓存数据_A", "本地缓存数据_B");

primaryRedisFlow
// 安全阀 2:捕获前面的异常,并“把整个身子切换到”备用流继续流淌
.onErrorResume(e -> {
System.err.println("警告:" + e.getMessage() + " 正在无缝切换到本地缓存...");
return backupLocalCacheFlow;
})
.subscribe(System.out::println);

Thread.sleep(2000);
}
1
2
3
4
5
Redis数据_1
Redis数据_2
本地缓存数据_A
本地缓存数据_B
警告:Redis 集群崩了! 正在无缝切换到本地缓存...


retry 或 retryWhen

顽强抵抗,原地无限/有限重试。有些网络抖动是暂时的,盲目降级很可惜。Reactor 提供了惊人的 retry 机制。它的底层原理是,一旦收到 onError 信号,自动重新订阅(resubscribe)最上游的数据源,把管道重新接通再跑一次。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
public static void main(String[] args) throws InterruptedException {
Flux.just("请求数据")
.map(data -> {
if (Math.random() > 0.3) { // 70% 概率网络抖动报错
System.out.println("❌ 尝试发起网络调用... 失败了!");
throw new RuntimeException("网络抖动异常");
}
return "🎉 成功拿到远程数据!";
})
// 安全阀 3:如果报错,自动重试 3 次,且每次重试之间间隔 200 毫秒(指数退避策略方案)
.retryWhen(reactor.util.retry.Retry.backoff(3, Duration.ofMillis(200)))
.subscribe(
System.out::println,
err -> System.err.println("😭 重试 3 次仍然失败,最终暴雷:" + err.getCause().getMessage())
);

Thread.sleep(2000); // 异步等待
}
1
2
3
❌ 尝试发起网络调用... 失败了!
❌ 尝试发起网络调用... 失败了!
🎉 成功拿到远程数据!

上述 过程,异常只是在事件循环中作为一个 Throwable 对象在方法入参(onError(t))里传递,没有任何一个物理线程因为这个异常而被迫终止或中断销毁。线程池里的 Worker 线程依然在优哉游哉地处理下一个事件,系统稳如泰山。


庞大的操作符家族

Project Reactor 的核心库 reactor-core 中包含了庞大的操作符家族。如果想完整、权威地查阅所有操作符,唯一的圣经就是其官方的 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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
public static void main(String[] args) {
// generate 接收两个主要部分:
// 1. 初始状态(这里提供了一个 Supplier,初始状态值为 1)
// 2. 状态转移函数(BiFunction:入参是当前状态和发射器,返回值是下一个新状态)
Flux<String> generateFlux = Flux.generate(
() -> 1, // 初始状态 state = 1
(state, sink) -> {
// 根据当前状态,吐出一个加工后的数据
sink.next("当前状态值: " + state + " -> 乘以3 = " + (state * 3));

// 终止条件控制:如果状态到了 9,拉响警报,主动关闭流水线
if (state >= 9) {
sink.complete();
}

// 返回下一个轮次的新状态(关键:state + 1 变成了下一轮的输入)
return state + 1;
}
);

// 订阅放水
generateFlux.subscribe(System.out::println);
}
1
2
3
4
5
6
7
8
9
当前状态值: 1 -> 乘以3 = 3
当前状态值: 2 -> 乘以3 = 6
当前状态值: 3 -> 乘以3 = 9
当前状态值: 4 -> 乘以3 = 12
当前状态值: 5 -> 乘以3 = 15
当前状态值: 6 -> 乘以3 = 18
当前状态值: 7 -> 乘以3 = 21
当前状态值: 8 -> 乘以3 = 24
当前状态值: 9 -> 乘以3 = 27

利用 generate 维护一个 “双变量状态(数组或元组)”,可以极其优雅地吐出一个完美的斐波那契响应式流。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
public static void main(String[] args) {

// 初始状态:一个包含两个元素的数组 [0, 1](代表斐波那契的前两项)
Flux<Long> fibonacciFlux = Flux.generate(
() -> new long[]{0L, 1L},
(state, sink) -> {
// 1. 吐出当前位置的斐波那契数
sink.next(state[0]);

// 2. 计算下一个状态值
long nextVal = state[0] + state[1];

// 3. 溢出保护:如果超过了 Long 的最大值就结束
if (nextVal < 0) {
sink.complete();
}

// 4. 状态转移:更新数组 [当前项变成下一项, 下一项变成刚才算出的和]
return new long[]{state[1], nextVal};
}
);

// 只拿前 10 个数字打印
fibonacciFlux.take(10).subscribe(num -> System.out.print(num + " "));
}
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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
import reactor.core.publisher.Flux;
import reactor.core.publisher.FluxSink;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;

public class FluxCreateDemo {

// 模拟传统的事件监听器接口
interface UserEventListener {
void onUserLogin(String username);
void onServerShutdown();
}

// 模拟一个事件源头(比如旧的 MQ 或者网关)
static class LegacyEventSource {
private final List<UserEventListener> listeners = new ArrayList<>();

public void register(UserEventListener listener) { this.listeners.add(listener); }

// 模拟突发事件产生
public void emitLogin(String name) {
for (UserEventListener l : listeners) l.onUserLogin(name);
}
public void emitShutdown() {
for (UserEventListener l : listeners) l.onServerShutdown();
}
}

public static void main(String[] args) throws InterruptedException, IOException {
LegacyEventSource legacySource = new LegacyEventSource();

// 🌟 核心魔法:使用 Flux.create 将老旧的事件源桥接到 Flux 世界
Flux<String> loginEventFlux = Flux.create(sink -> {

// 1. 在传统的监听器中,把收到的数据无缝丢进响应式 sink(发射器)
legacySource.register(new UserEventListener() {
@Override
public void onUserLogin(String username) {
// 核心动作:来一个,我吐一个,甚至可以连着吐好几个,完全没有 generate 的次数限制
sink.next("👤 [事件捕获] 用户登录: " + username);
}

@Override
public void onServerShutdown() {
sink.complete();
System.out.println("收到结束信号,关闭响应式管道");
}
});

// 2. 联动注销:当下游主动 dispose 取消订阅时,顺便把老旧监听器注销,防止内存泄漏
sink.onCancel(() -> System.out.println("🛑 下游断开,解除老旧监听器的绑定关系"));

}, FluxSink.OverflowStrategy.BUFFER); // 🌟 指定背压策略:如果数据积压,先用内存 Buffer 接着

// 订阅放水,开始监控
loginEventFlux.subscribe(System.out::println);

// 模拟外部世界在长达几秒的时间内,零星或者密集地产生事件(可以是任何线程触发)
legacySource.emitLogin("张三");
Thread.sleep(200);
legacySource.emitLogin("李四");
legacySource.emitLogin("王五"); // 👈 可以在单次瞬间连着吐两个数据,完全自由
Thread.sleep(200);

// 结束服务器,触发 complete
legacySource.emitShutdown();
}
}
1
2
3
4
👤 [事件捕获] 用户登录: 张三
👤 [事件捕获] 用户登录: 李四
👤 [事件捕获] 用户登录: 王五
收到结束信号,关闭响应式管道


自定义消费者

虽然普通的 .subscribe(data -> { … }) 传入一个 Consumer 很方便,但它是一个 “全自动” 的水龙头——一旦打开,上游会以排山倒海之势把所有数据一口气推下来(默认请求 Long.MAX_VALUE)。如果你想实现“手控水阀”(比如:下游有空时才分批要数据,或者动态拒绝数据),你就必须通过继承 BaseSubscriber 来自己造一个全功能的“水龙头”。在 Reactor 中,如果你想完全控制消费者的底层行为(比如手动精确控制背压、监听各种生命周期钩子信号),BaseSubscriber 是官方最推荐、最正宗的唯一首选基类。

你可能会想,既然有规范定义的 Subscriber 接口,我自己写个类 implements Subscriber 不行吗?官方不建议这样做。因为直接实现接口,你需要自己去处理复杂的线程安全问题、Subscription 状态维护、以及极其容易发生内存泄漏的并发规则。 而 BaseSubscriber 已经把这些底层的“脏活累活”全部封装好了,并提供了安全的保护机制。

我们来看一个典型的生产案例。上游流水线疯狂流出数据,自定义消费者因为处理能力有限,必须实行严格的 “限流分批” 策略:第一批只要 3 个,每处理完 1 个,才允许上游再送 1 个(动态维持小缓冲区)。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
import org.reactivestreams.Subscription;
import reactor.core.publisher.BaseSubscriber;
import reactor.core.publisher.Flux;

public class CustomSubscriberDemo {
public static void main(String[] args) {
Flux<String> fastSource = Flux.just("元素1", "元素2", "元素3", "元素4", "元素5");
fastSource.subscribe(new BatchControlSubscriber()); // 使用我们自定义的“手控水龙头”进行订阅
}
}

/**
* 继承 BaseSubscriber 自定义的高级限流消费者
*/
class BatchControlSubscriber extends BaseSubscriber<String> {

// 1. 当订阅关系刚建立时的“空间钩子”
@Override
protected void hookOnSubscribe(Subscription subscription) {
System.out.println("🔌 管道接通!");
// 核心动作:跟上游说,我刚开张,第一批先给我放 3 滴水,别给多了!
request(3);
}

// 2. 每当真实数据(onNext 信号)流过来时的加工钩子
@Override
protected void hookOnNext(String value) {
System.out.println(" [消费者] 正在精心加工: " + value);

// 模拟复杂的业务消耗...

// 核心控速:加工完这 1 个之后,再次向上游申请追加 1 个名额
System.out.println(" --> 报销完毕,向上游追加请求 1 个名额");
request(1);
}

// 3. 当流正常流完时的钩子
@Override
protected void hookOnComplete() {
System.out.println("🎉 上游水放完了,流水线完美收工!");
}

// 4. 当中途发生暴雷异常时的钩子
@Override
protected void hookOnError(Throwable throwable) {
System.err.println("🚨 消费者拉响警报,发现异常: " + throwable.getMessage());
}

// 5. 无论是正常结束、报错,还是中途取消,最后一定会执行的清理钩子
@Override
protected void hookOnCancel() {
System.out.println("🛑 下游主动关停了水阀,做最后的资源清理...");
}
}
1
2
3
4
5
6
7
8
9
10
11
12
🔌 管道接通!
[消费者] 正在精心加工: 元素1
--> 报销完毕,向上游追加请求 1 个名额
[消费者] 正在精心加工: 元素2
--> 报销完毕,向上游追加请求 1 个名额
[消费者] 正在精心加工: 元素3
--> 报销完毕,向上游追加请求 1 个名额
[消费者] 正在精心加工: 元素4
--> 报销完毕,向上游追加请求 1 个名额
[消费者] 正在精心加工: 元素5
--> 报销完毕,向上游追加请求 1 个名额
🎉 上游水放完了,流水线完美收工!


Disposable 接口

Disposable 的使用

Disposable 是你调用 .subscribe() 拧开水龙头后,手里拿到的那个 “遥控关阀开关”。当你启动了一个无限流(比如每秒发射一次的定时器,或者一个大文件的读取任务),如果下游在某个时刻不想听了、或者请求已经超时,就可以直接调用 disposable.dispose()。此时下游会向上一路逆流发送一个 cancel(取消)信号通知上游立刻停止生产,从而彻底释放物理线程、网络 Socket 或文件句柄,防止内存泄漏。

1
2
3
4
5
public interface Disposable {
void dispose(); // 🔴 咔哒!强制切断数据流
boolean isDisposed(); // 检查这个开关是否已经被按下了
// ...
}

我们来看一个最直观的例子。我们启动一个每隔 500 毫秒就疯狂报时的闹钟流,并在 2 秒后用 Disposable 手动取消它。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
public static void main(String[] args) throws InterruptedException {

// 1. 创建一个无限流:每 500ms 报时一次
Flux<Long> clockFlux = Flux.interval(Duration.ofMillis(500))
.doOnCancel(() -> System.out.println("🛑 【通知】收到下游的取消信号,上游定时器已安全关闭。"));

// 2. 订阅并拿到“遥控关阀开关” Disposable
Disposable subscriptionControl = clockFlux.subscribe(
tick -> System.out.println("⏰ 闹钟报时: 滴答 " + tick)
);

// 主线程睡 2000 毫秒(2 秒),让闹钟走一会儿
Thread.sleep(2000);

// 3. 核心动作:下游觉得吵了,按下开关,强制关停!
System.out.println("⏳ 2秒时间到,下游发起强制关停...");
subscriptionControl.dispose();

// 再睡 1 秒,观察闹钟还会不会继续响
Thread.sleep(1000);
System.out.println("🎉 程序结束,无任何多余线程驻留。");
}
1
2
3
4
5
6
7
⏰ 闹钟报时: 滴答 0
⏰ 闹钟报时: 滴答 1
⏰ 闹钟报时: 滴答 2
⏳ 2秒时间到,下游发起强制关停...
⏰ 闹钟报时: 滴答 3
🛑 【通知】收到下游的取消信号,上游定时器已安全关闭。
🎉 程序结束,无任何多余线程驻留。


Disposables 复合开关舱

在写响应式项目时,一个用户可能会同时打开好几个数据订阅(比如:接收私聊流、接收群聊流、接收系统通知流)。当这个用户下线(断开 WebSocket)时,我们必须要同时关闭这几百个流。如果一个一个去调用 dispose() 会非常痛苦。Reactor 提供了一个完美的聚合工具类:

1
2
3
4
5
6
7
8
9
10
11
12
13
import reactor.core.Disposable;
import reactor.core.Disposables;

// 1. 创建一个“复合开关舱”
Disposable.Composite compositeDisposable = Disposables.composite();

// 2. 把多个订阅的遥控器全部丢进舱内
compositeDisposable.add(chatFlux.subscribe());
compositeDisposable.add(groupFlux.subscribe());
compositeDisposable.add(noticeFlux.subscribe());

// 3. 当用户下线时,一键引爆,全部销毁!
compositeDisposable.dispose();


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
2
3
4
5
6
7
8
Mono.just("Hello")
.contextWrite(Context.of("KEY", "VALUE-A")) // ❌ 写在了读取的上面
.flatMap(data -> Mono.deferContextual(ctx -> {
// 报错!因为逆流而上时,信号根本到不了上面那行 contextWrite
String val = ctx.get("KEY");
return Mono.just(data + " " + val);
}))
.subscribe();

正确示范:把注入写在读取的“下游”(后面)

1
2
3
4
5
6
7
Mono.just("Hello")
.flatMap(data -> Mono.deferContextual(ctx -> {
String val = ctx.get("KEY"); // 完美拿到!因为下游的 contextWrite 逆流冲刷了这一层
return Mono.just(data + " " + val);
}))
.contextWrite(Context.of("KEY", "VALUE-A")) // 🌟 写在后面(下游)
.subscribe(System.out::println);

简单案例:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
public static void main(String[] args) {
// 1. 声明一条业务流水线(模拟 Service 层 + DAO 层)
Mono<String> businessPipeline = Mono.deferContextual(ctx -> {
// 🌟 核心:通过 deferContextual 现场去捞当前上下文里的 TraceId
String traceId = ctx.getOrDefault("TRACE_ID", "unknown-id");
System.out.println("[DAO层] [" + traceId + "] 正在执行数据库查询,当前线程: " + Thread.currentThread().getName());
return Mono.just("[DAO层结果] [" + traceId + "] 成功拿到用户数据");
})
.map(data -> {
// 模拟中间转换
return "[Service层加工] " + data;
});

// 2. 模拟 Controller 层:接收到请求,并给该请求打上唯一的 TRACE_ID
System.out.println("用户的 HTTP 请求进入 Controller...");

businessPipeline
// 🌟 核心动作:向下游/在尾部注入 Context!
.contextWrite(Context.of("TRACE_ID", "8888-aaaa-9999-bbbb"))
.subscribe(result -> {
System.out.println("[前端收到响应] " + result);
});
}
1
2
3
用户的 HTTP 请求进入 Controller...
[DAO层] [8888-aaaa-9999-bbbb] 正在执行数据库查询,当前线程: main
[前端收到响应] [Service层加工] [DAO层结果] [8888-aaaa-9999-bbbb] 成功拿到用户数据


使用新版 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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
import reactor.core.publisher.Mono;
import reactor.netty.http.server.HttpServer;

public class PureReactorServer {

public static void main(String[] args) {
System.out.println("⏳ 正在利用纯原生 Reactor API 编排服务器...");

// 1. 创建并配置一个原生的 Reactor Http 服务器
reactor.netty.DisposableServer server = HttpServer.create()
.port(8088) // 监听 8088 端口
.route(routes -> routes
// 🧱 路由一:最简单的 GET 请求响应
.get("/hello", (request, response) ->
response.sendString(Mono.just("Hello! This is a pure Reactor Server without Spring!"))
)

// 🌊 路由二:处理 POST 请求,读取请求体(Request Body)并进行响应式转换
// 逻辑:把用户发过来的文本,转换成大写后原路吐回去
.post("/echo", (request, response) ->
response.sendString(
request.receive() // 1.接收字节流
.asString() // 2.转化为字符串 Flux 流
.map(String::toUpperCase) // 3️⃣3.原生操作符:转换为大写
.doOnNext(body -> System.out.println("📝 [服务器收到并处理POST数据]: " + body))
)
)
)
.bindNow(); // ⚡ 瞬间绑定并启动服务器

System.out.println("🚀 原生 Reactor Netty 服务器已成功在端口 8088 启航!");

// 2. 阻塞主线程以使服务器保持存活(生产级标准写法)
server.onDispose().block();
}
}

直接右键运行上方的 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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
package com.zdemo.scloud;

import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;
import reactor.netty.DisposableServer;
import reactor.netty.http.server.HttpServer;
import reactor.netty.http.server.HttpServerResponse;

import java.io.FileNotFoundException;
import java.time.Duration;

public class PureReactorServer {

public static void main(String[] args) {
System.out.println("⏳ ChatGPT 模拟服务器正在无声启动...");

DisposableServer server = HttpServer.create()
.port(8088)
.route(routes -> routes
// 🌟 核心:ChatGPT 流式对话接口
.post("/api/chat", (request, response) -> {

// 1️⃣ 设置标准的 SSE 响应头,允许跨域(方便前后端分离测试)
setSseHeaders(response);

// 2️⃣ 接收前端传过来的用户问题
return response.sendString(
request.receive()
.asString()
.flatMap(userPrompt -> {
System.out.println("🤖 [收到用户提问]: " + userPrompt);

// 3️⃣ 模拟 AI 生成的深度回答文本
String aiAnswer = "你好!我是基于纯原生 Project Reactor 构建的 AI 助手。感谢你对 koohub(koo hub, cool code.)的关注!"
+ "你正在体验的是 ChatGPT 同款的 SSE(Server-Sent Events)流式推流技术。"
+ "在整个输出过程中,底层的 Netty 线程没有发生任何一点阻塞,这就是响应式的终极魅力!";

// 4️⃣ 将一整段文本切碎成单个字符数组
String[] charArray = aiAnswer.split("");

// 5️⃣ 核心响应式编排:利用高频时钟把字符吐出去
return Flux.interval(Duration.ofMillis(50)) // 每 50ms 发出一个滴答信号
.take(charArray.length) // 答完即止,有多少个字就走多少次
.map(index -> {
String currentChar = charArray[index.intValue()];
// 🌟 严格遵守标准 SSE 协议格式: data: 内容\n\n
return "data: " + currentChar + "\n\n";
});
})
);
})
// 额外提供一个静态页面路由,方便前端直接访问
.get("/", (request, response) -> {
Mono<byte[]> asyncFileMono = Mono.fromCallable(() -> {
try (var inputStream = PureReactorServer.class.getClassLoader().getResourceAsStream("index.html")) {
if (inputStream == null) {
throw new FileNotFoundException();
}
return inputStream.readAllBytes();
}
}).subscribeOn(Schedulers.boundedElastic());

return asyncFileMono
.flatMap(bytes -> response.sendByteArray(Mono.just(bytes)).then())
.onErrorResume(FileNotFoundException.class, err ->
response.status(404)
.sendString(Mono.just("Error 404: index.html not found!"))
.then()
);
})
)
.bindNow();

System.out.println("🚀 ChatGPT 模拟微服务已在端口 8088 扬帆起航!");
server.onDispose().block();
}

/**
* 设置标准的 SSE 和跨域响应头
*/
private static void setSseHeaders(HttpServerResponse response) {
response.header("Content-Type", "text/event-stream;charset=UTF-8") // 🌟 必须:指定为事件流
.header("Cache-Control", "no-cache") // 🌟 必须:禁用缓存,确保实时
.header("Connection", "keep-alive") // 🌟 必须:保持长连接
.header("Access-Control-Allow-Origin", "*") // 允许跨域
.header("Access-Control-Allow-Headers", "*");
}
}

在项目 resources 目录,创建一个名为 index.html 的文件。传统的 EventSource 只支持 GET 请求。为了配合 AI 业务中常见的 POST 提问,前端使用现代浏览器的 fetch API 和 ReadableStream(可读流) 来动态拦截并渲染后端的响应式流。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
<!DOCTYPE html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8">
<title>KooHub AI - 纯原生 Reactor 体验舱</title>
<style>
body { font-family: 'Segoe UI', Tahoma, Geneva, Verdana, sans-serif; background: #f4f7f6; padding: 40px; }
.chat-container { max-width: 700px; margin: 0 auto; background: white; padding: 30px; border-radius: 12px; box-shadow: 0 4px 15px rgba(0,0,0,0.05); }
h2 { color: #2c3e50; border-bottom: 2px solid #ecf0f1; padding-bottom: 15px; margin-top: 0; }
.input-box { display: flex; margin-bottom: 20px; }
input { flex: 1; padding: 12px; border: 1px solid #ddd; border-radius: 6px; font-size: 16px; outline: none; }
button { background: #1abc9c; color: white; border: none; padding: 0 24px; margin-left: 10px; border-radius: 6px; cursor: pointer; font-size: 16px; transition: background 0.3s; }
button:hover { background: #16a085; }
.output-zone { background: #2c3e50; color: #ecf0f1; padding: 20px; border-radius: 8px; min-height: 150px; font-size: 16px; line-height: 1.6; white-space: pre-wrap; word-break: break-all; }
.cursor { display: inline-block; width: 2px; height: 18px; background: #1abc9c; animation: blink 1s infinite; vertical-align: middle; margin-left: 2px; }
@keyframes blink { 0%, 100% { opacity: 0; } 50% { opacity: 1; } }
</style>
</head>
<body>

<div class="chat-container">
<h2>🤖 KooHub AI 流式对话舱</h2>
<div class="input-box">
<input type="text" id="promptInput" value="请介绍一下响应式编程的核心优势" placeholder="向 AI 提问...">
<button onclick="sendStreamRequest()">发射提问</button>
</div>
<div class="output-zone" id="outputZone">等待提问中...</div>
</div>

<script>
async function sendStreamRequest() {
const inputEl = document.getElementById('promptInput');
const outputEl = document.getElementById('outputZone');
const prompt = inputEl.value;

if (!prompt) return;

// 清空舞台,放置打字机光标
outputEl.innerHTML = '<span class="cursor"></span>';

try {
// 1 发起异步 POST 请求
const response = await fetch('http://localhost:8088/api/chat', {
method: 'POST',
body: prompt,
headers: { 'Content-Type': 'text/plain' }
});

// 2 拿到底层的非阻塞字节流读取器 (Reader)
const reader = response.body.getReader();
const decoder = new TextDecoder('utf-8');

outputEl.innerHTML = ""; // 正式开始接收,移除光标前的空白

// 3 循环死抠管道里的每一滴水(字节)
while (true) {
const { value, done } = await reader.read();
if (done) break; // 水流抽干,安全退出

// 4 解码字节块为文本
const chunk = decoder.decode(value, { stream: true });

// 5 按照标准的 SSE 协议解析数据 (过滤掉前面的 "data: " 和后面的换行)
const lines = chunk.split('\n');
for (const line of lines) {
if (line.startsWith('data: ')) {
const actualChar = line.replace('data: ', '');
// 🌟 动态追加到前端界面,形成流式打字机视觉效果
outputEl.innerHTML += actualChar;
}
}
}
// 结束渲染后,加上完美收尾标识
outputEl.innerHTML += ' <b style="color:#1abc9c;">✓</b>';

} catch (error) {
outputEl.innerHTML = "❌ 发生致命通信错误: " + error.message;
}
}
</script>
</body>
</html>

联调和效果展示:

  • 启动后端:运行 PureReactorServer 的 main 方法。
  • 打开前端:在浏览器直接输入并打开地址:http://localhost:8088/ 或者是双击本地的 index.html。
  • 点击“发射提问”:你会看到下方的黑色控制台区域,完全没有长时间的空白等待,而是像 ChatGPT 官方一样,文字以每秒 20 个字的高频速度,如丝般顺滑地“吐”在屏幕上!