随着用户规模与数据吞吐量的激增,传统的同步阻塞模型越来越难以满足高并发、低延迟的需求。线程池虽然能在一定程度上缓解问题,但有限线程数和上下文切换开销仍是瓶颈。响应式编程(Reactive Programming) 正是在这一背景下进入主流视野,而 Spring 生态中的 Project Reactor 则成为 Java 平台上响应式开发的事实标准。
25.1.1 响应式编程的核心思想
响应式编程并非新鲜概念,它是一种以数据流和变化传播为中心的异步编程范式。与传统命令式编程“一步一步执行”不同,响应式编程让你定义一条数据处理的流水线,数据作为“事件”在流水线中流动,触发各个处理环节。
其核心思想可以归纳为以下几点:
1. 异步非阻塞
同步阻塞模型中,每个请求都会占用一个线程,该线程在等待 I/O(数据库查询、远程调用)时完全挂起,无法处理其他工作。响应式编程则基于事件驱动,线程发起一个异步操作后立即返回,当数据就绪或发生错误时,由回调或观察者来处理结果。线程数量大大减少,线程资源得到更高效的利用。
2. 数据流(Stream)与事件驱动
一切皆流。响应式编程把任何事物(数据库返回的多条记录、鼠标点击事件、消息队列的消息、HTTP 请求体)都视为一个事件序列。你不再去“拉取”数据,而是订阅一个数据源,当数据到达时,流水线自动触发相应的处理函数。这种思维方式的转变,使处理逻辑变成了对数据流的声明式变换。
3. 声明式组合与操作符
你不需要手动编写复杂的线程调度和同步控制代码,而是通过丰富的操作符(map、filter、flatMap、merge、zip 等)对数据流进行转换、过滤、组合。代码变得高度可读和可维护,形成一条清晰的“装配线”。
4. 背压(Backpressure)机制
这是响应式系统中至关重要的概念。当生产者生产数据的速度超过消费者的处理能力时,系统必须有反馈机制来协调速度,否则会耗尽资源导致崩溃。背压就是消费者向上游通知“我现在能处理多少”的能力。Reactor 底层支持背压传播,确保整个流式链路上的速度匹配。
25.1.2 Reactor 基础:Mono 与 Flux
Project Reactor 是 Spring 5 引入的响应式核心库,它实现了 Reactive Streams 规范,并提供两种核心类型:Mono 和 Flux。它们代表不同基数的异步数据流:
- Mono<T>:表示 0 到 1 个元素的异步序列。适合单一结果场景,例如根据 ID 查询单条记录、发送一条消息、或者一个可以完成或失败的空操作。
- Flux<T>:表示 0 到 N 个元素的异步序列。适用于多条记录、事件流、实时数据等场景。
可以借助简单记忆:Mono 是“最多一个”,Flux 是“多个”。
两者都是惰性的:只有当你调用 subscribe() 时,数据流才会启动。没有订阅,流水线只是一段描述,不会有任何实际动作。
创建 Mono 与 Flux
// 创建包含单值的 Mono
Mono<String> mono = Mono.just("Hello, Reactive World");
// 创建一个空的 Mono
Mono<Object> empty = Mono.empty();
// 从 Callable 创建(可能抛异常)
Mono<String> fromCallable = Mono.fromCallable(() ->
db.findById(1L)); // 假设 db 为某个接口
// 创建包含多个元素的 Flux
Flux<String> flux = Flux.just("Spring", "Reactor", "WebFlux");
// 从 Iterable 创建
Flux<Integer> range = Flux.range(1, 5); // 1,2,3,4,5
// 从 Stream 创建
Flux<String> fromStream = Flux.fromStream(Stream.of("a", "b", "c"));
订阅与消费
Flux<String> cities = Flux.just("Beijing", "Shanghai", "Shenzhen");
cities.subscribe(); // 启动流但不处理数据(仅触发副作用)
cities.subscribe(
city -> System.out.println("City: " + city), // 数据消费
error -> System.err.println("Error: " + error), // 错误处理
() -> System.out.println("All cities processed") // 完成回调
);
只有调用 subscribe,数据才开始流动,这正是惰求值的体现。
25.1.3 核心操作符:像流一样思考
Reactor 的强大在于提供了上百个操作符,让你以函数式风格操作数据流。下面介绍几种最常用、最实用的操作符。
1. 转换:map 与 flatMap
map对每个元素同步转换,返回普通值:
Flux<String> upper = Flux.just("reactive", "programming")
.map(String::toUpperCase); // "REACTIVE", "PROGRAMMING"
flatMap将每个元素映射为一个异步的Publisher(Mono 或 Flux),然后将这些内部流扁平化合并成一个大的流。特别适合需要异步调用的场景,例如根据 ID 调用远程服务:
Flux<Order> orders = Flux.just("order1", "order2")
.flatMap(orderId -> orderService.findById(orderId)); // 返回 Mono<Order>
注意:flatMap 不保证顺序,内部请求是并发执行的。如果需要保持顺序,可使用 concatMap。
2. 过滤:filter
Flux<Integer> even = Flux.range(1, 10)
.filter(i -> i % 2 == 0); // 2,4,6,8,10
3. 组合:zip 与 merge
zip将多个流“拉链式”配对组合,等所有流都有元素时才生成组合结果:
Flux<String> firstNames = Flux.just("John", "Jane");
Flux<String> lastNames = Flux.just("Doe", "Smith");
Flux<String> fullNames = Flux.zip(firstNames, lastNames,
(f, l) -> f + " " + l); // "John Doe", "Jane Smith"
merge将多个流合并成一个流,元素到达的顺序就是实际到达的顺序,不要求配对。
Flux<String> merged = Flux.merge(fluxA, fluxB);
4. 错误处理
Reactor 提供多种优雅的错误处理方式,避免回调地狱。
Flux.just("1", "2", "three")
.map(Integer::parseInt) // "three" 将导致异常
.onErrorReturn(-1) // 错误时返回默认值 -1
.subscribe(System.out::println); // 1,2,-1
其他常用操作符:onErrorResume(回退到备用流)、doOnError(记录日志但不中断流)、retry(重试)。
5. 背压控制(体现响应式核心能力)
消费者可以通过 request(n) 控制需求,但更常用的是利用操作符间接控制:
Flux.range(1, 1000)
.limitRate(10) // 每次请求最多10个元素
.subscribe(data -> process(data));
在 Spring WebFlux 中,底层网络层会自动处理背压,开发者通常无需手动干预,但理解背压机制对排查性能问题至关重要。
25.1.4 Reactor 与命令式编程的融合实践
React 没有强制你完全放弃命令式代码,通过 block() 和 blockOptional() 可以在必要时同步等待结果,用于与遗留代码交互或测试:
Mono<String> result = Mono.just("Hello").delayElement(Duration.ofMillis(100));
String value = result.block(); // 阻塞等待,仅用于测试或边缘场景!
但在生产代码中应该避免 block(),因为这会丧失非阻塞的优势。测试时,更好的方式是使用 StepVerifier:
StepVerifier.create(Flux.range(1, 3))
.expectNext(1, 2, 3)
.verifyComplete();
25.1.5 与 Spring 生态的深度整合
Reactor 本身是一个通用库,但它在 Spring 生态中的应用几乎无处不在:
- Spring WebFlux 的控制器可以直接返回
Mono/Flux,框架负责订阅并将数据序列化输出。 - Spring Data Reactive(如 Reactive MongoDB、R2DBC)提供响应式的 Repository,方法返回
Mono或Flux。 - Spring Cloud Gateway 使用 Reactor 构建非阻塞网关路由。
- Spring Security Reactive 提供响应式的安全过滤链。
这意味着,掌握 Reactor 基础后,整个响应式 Spring 技术栈的大门就将打开。你编写的代码不再是一段段串行指令,而是一张张数据流拓扑图,JVM 以极少的线程吞吐海量请求,这正是现代高性能应用所需要的能力。
本节梳理了响应式编程的核心思想和 Reactor 的核心类型及操作符。下一节,我们将在此基础上踏入 Spring WebFlux 的大门,亲手构建第一个非阻塞的 RESTful 服务。