响应式编程基础
响应式编程(Reactive Programming)是一种基于数据流(Data Stream)和变化传播的编程范式。在 Spring WebFlux 体系中,它使得应用能够以声明式的方式处理异步数据流,用少量线程支撑高并发场景。
本文从 Reactive Streams 规范出发,深入 Flux/Mono 操作符、调度器、背压机制以及实战中的 R2DBC 连接池差异分析,帮助读者建立完整的响应式编程知识体系。
1. Reactive Streams 规范
1.1 背景与目标
Reactive Streams 是由 Netflix、Lightbend、Pivotal 等公司联合制定的异步流处理标准,定义了一套带背压的异步数据流处理接口。其核心目标在于:
- 异步非阻塞:绝不阻塞调用线程
- 背压驱动:订阅者能主动控制接收速率
- 可组合性:通过操作符链式组合
- 低延迟、高吞吐:适用于网络 I/O 密集型应用
1.2 四大核心接口
Reactive Streams 规范定义四个核心接口:
| 接口 | 角色 | 说明 |
|---|---|---|
Publisher<T> | 数据发布者 | 生产数据,接受 Subscriber 订阅 |
Subscriber<T> | 数据订阅者 | 消费数据,处理 onNext/onError/onComplete |
Subscription | 订阅契约 | 连接 Publisher 和 Subscriber,提供 request(n) 与 cancel() |
Processor<T,R> | 处理节点 | 同时是 Publisher 和 Subscriber |
// 规范的典型交互流程
Publisher --subscribe--> Subscriber
Subscriber --onSubscribe(Subscription)--> Publisher
Subscriber --request(n)--> Publisher // 请求 n 个元素
Publisher --onNext(T)--> Subscriber // 逐个推送
...循环直到...
Publisher --onComplete()--> Subscriber // 正常结束
// 或
Publisher --onError(Throwable)--> Subscriber // 异常结束1.3 Spring WebFlux 中的实现
Spring WebFlux 底层依赖 Project Reactor,提供了两个核心发布者:
// Flux 与 Mono 的语义对比
Flux<T> —— 0..N 个元素 —— 类比 Collection<T>
Mono<T> —— 0..1 个元素 —— 类比 Optional<T>
Mono<Void> —— 只关注完成信号 —— 类比 Runnable2. Flux / Mono 核心操作符
2.1 map —— 同步转换
map 对每个元素执行同步转换函数,一对一的元素映射。
Flux<Integer> source = Flux.just(1, 2, 3, 4, 5);
Flux<String> result = source.map(i -> "Number: " + i);
// 输出: Number: 1, Number: 2, Number: 3, Number: 4, Number: 5注意:map 内不宜执行阻塞操作或耗时调用,否则会阻塞整个流水线线程。
2.2 flatMap —— 异步展平
flatMap 将每个元素映射为一个 Publisher,然后合并展平成一个新的 Flux。元素间并发执行,结果顺序可能乱序。
Flux<String> words = Flux.just("hello", "world");
Flux<String> letters = words.flatMap(w -> Flux.fromArray(w.split("")));
// 输出: h, e, l, l, o, w, o, r, l, d (顺序不确定)2.3 filter —— 过滤
filter 根据断言(Predicate)保留满足条件的元素。
Flux<Integer> source = Flux.range(1, 10);
Flux<Integer> even = source.filter(i -> i % 2 == 0);
// 输出: 2, 4, 6, 8, 102.4 zip —— 一对一合并
zip 将多个 Publisher 的对应元素按一对一模式合并,任一 Publisher 完成后则结束。
Flux<String> names = Flux.just("A", "B", "C");
Flux<Integer> scores = Flux.just(90, 85, 95);
Flux<String> combined = Flux.zip(names, scores,
(name, score) -> name + ":" + score);
// 输出: A:90, B:85, C:952.5 merge —— 交错合并
merge 将多个 Publisher 交错合并,哪个 Publisher 先来数据就优先发出,不保证顺序。
Flux<Long> fast = Flux.interval(Duration.ofMillis(100)).take(5);
Flux<Long> slow = Flux.interval(Duration.ofMillis(200)).take(3);
Flux<Long> merged = Flux.merge(fast, slow);
// 输出(示例): 0(fast), 0(slow), 1(fast), 2(fast), 1(slow), 3(fast), 4(fast), 2(slow)2.6 concat —— 顺序连接
concat 按顺序串联多个 Publisher,前一个完成后才订阅下一个,保证顺序。
Flux<String> first = Flux.just("A1", "A2");
Flux<String> second = Flux.just("B1", "B2");
Flux<String> concatenated = Flux.concat(first, second);
// 输出: A1, A2, B1, B2(严格有序)2.7 retry —— 重试
当流发生错误时,retry 会重新订阅上游 Publisher,从头开始重新发射所有元素。
Flux<String> unstable = Flux.just("ok", "ok")
.concatWith(Flux.error(new RuntimeException("DB error")))
.concatWith(Flux.just("never"));
unstable.retry(2)
.subscribe(
data -> System.out.println("收到: " + data),
err -> System.err.println("最终失败: " + err.getMessage())
);
// 原始 + 2 次重试,共执行 3 次2.8 onErrorResume —— 错误回退
当发生错误时,onErrorResume 提供一个备用 Publisher 来替代错误信号,使流能继续。
Flux<String> safe = Flux.just("data1", "data2")
.concatWith(Flux.error(new RuntimeException("failed")))
.onErrorResume(e -> {
System.err.println("捕获错误: " + e.getMessage());
return Flux.just("fallback1", "fallback2");
});
// 输出: data1, data2, fallback1, fallback23. subscribe 方式与自定义 Subscriber
3.1 subscribe 的多种重载
Reactor 提供了多级粒度的 subscribe 方法:
// 方式1:完全忽略信号
Flux.just(1, 2, 3).subscribe();
// 方式2:只处理 onNext
Flux.just(1, 2, 3).subscribe(System.out::println);
// 方式3:处理 onNext + onError
Flux.just(1, 2, 3)
.subscribe(System.out::println, Throwable::printStackTrace);
// 方式4:处理 onNext + onError + onComplete
Flux.just(1, 2, 3)
.subscribe(
System.out::println,
Throwable::printStackTrace,
() -> System.out.println("完成")
);
// 方式5:全参数 + Subscription 消费
Flux.just(1, 2, 3)
.subscribe(
System.out::println,
Throwable::printStackTrace,
() -> System.out.println("完成"),
subscription -> subscription.request(Long.MAX_VALUE)
);3.2 自定义 Subscriber —— 细粒度背压控制
通过实现 org.reactivestreams.Subscriber 接口,可以精确控制每次请求的元素数量。
public class MySubscriber<T> implements Subscriber<T> {
private Subscription subscription;
private static final int REQUEST_SIZE = 2;
@Override
public void onSubscribe(Subscription s) {
this.subscription = s;
System.out.println("订阅成功,首次请求 " + REQUEST_SIZE + " 个元素");
s.request(REQUEST_SIZE);
}
@Override
public void onNext(T item) {
System.out.println("收到: " + item);
subscription.request(1); // 每次处理完再请求一个(动态背压)
}
@Override
public void onError(Throwable t) {
System.err.println("发生错误: " + t.getMessage());
}
@Override
public void onComplete() {
System.out.println("所有元素处理完毕");
}
}
Flux.range(1, 10).subscribe(new MySubscriber<>());3.3 BaseSubscriber 便捷方式
Reactor 提供了 BaseSubscriber 抽象类,使用时只需覆盖关心的回调:
Flux.range(1, 100)
.subscribe(new BaseSubscriber<Integer>() {
@Override
protected void hookOnSubscribe(Subscription subscription) {
request(5); // 初始请求 5 个
}
@Override
protected void hookOnNext(Integer value) {
System.out.println("处理: " + value);
if (value % 10 == 0) {
request(5); // 每 10 个元素追加请求
}
}
@Override
protected void hookOnComplete() {
System.out.println("完成");
}
@Override
protected void hookOnError(Throwable throwable) {
System.err.println("错误: " + throwable);
}
});4. Scheduler 调度器与线程模型
4.1 调度器的作用
默认情况下,操作符执行在订阅线程上。Scheduler 用于控制响应式链在哪个线程池上执行,实现异步边界。
4.2 常用调度器
| 调度器 | 线程模型 | 适用场景 |
|---|---|---|
Schedulers.immediate() | 当前线程执行 | 默认行为,测试 |
Schedulers.single() | 单一可复用线程 | 轻量级异步任务 |
Schedulers.boundedElastic() | 有界弹性线程池(默认 10 × CPU 核数) | 阻塞 I/O |
Schedulers.parallel() | 固定大小(CPU 核数)的工作窃取池 | CPU 密集型计算 |
Schedulers.newBoundedElastic(...) | 自定义有界弹性线程池 | 自定义隔离 |
4.3 切换调度器的操作符
// subscribeOn —— 影响订阅触发的线程(即数据发射的源头线程)
Flux.just("A", "B", "C")
.subscribeOn(Schedulers.boundedElastic())
.map(String::toLowerCase)
.subscribe(System.out::println);
// publishOn —— 影响后续操作符执行的线程
Flux.just("A", "B", "C")
.publishOn(Schedulers.parallel())
.map(s -> {
System.out.println("处理线程: " + Thread.currentThread().getName());
return s.toLowerCase();
})
.subscribeOn(Schedulers.single())
.subscribe(System.out::println);关键区别:
subscribeOn:改变订阅源头的线程,链中最靠近源的那个生效publishOn:改变下游操作符的执行线程,可多次使用形成多段线程切换
// 线程切换示意图
source.publishOn(s1) → op1 → publishOn(s2) → op2
// ↑ ↑
// op1 在 s1 执行 op2 在 s2 执行
// source 和 subscribe 在调用线程执行4.4 ParallelFlux —— CPU 密集型并行
Flux.range(1, 1_000_000)
.parallel(4)
.runOn(Schedulers.parallel())
.map(i -> expensiveCompute(i))
.sequential()
.subscribe(System.out::println);4.5 自定义 Scheduler 与线程饥饿规避
// 自定义 boundedElastic 线程池
Scheduler customElastic = Schedulers.newBoundedElastic(
50, 1000, "my-biz-pool", 60, false);
Flux.range(1, 100)
.flatMap(i ->
Mono.fromCallable(() -> ioOperation(i))
.subscribeOn(customElastic)
)
.subscribe();
// 应用关闭时:customElastic.dispose();线程饥饿:在响应式链路中将阻塞操作(如 JDBC)放在 parallel() 调度器上会阻塞 event-loop 线程,应始终用 boundedElastic() 隔离阻塞调用。
5. 背压(Backpressure)机制
5.1 什么是背压
背压是指**下游订阅者向上游发布者反馈"我处理不过来了,请慢一点"**的能力。它防止生产者速度超过消费者速度导致内存溢出或系统崩溃。
5.2 背压协议
背压通过 Subscription.request(n) 方法实现:
Flux.range(1, Integer.MAX_VALUE)
.subscribe(new BaseSubscriber<Integer>() {
@Override
protected void hookOnSubscribe(Subscription subscription) {
subscription.request(1); // 一次只请求 1 个
}
@Override
protected void hookOnNext(Integer value) {
process(value);
request(1); // 处理完后请求下一个
}
});5.3 Reactor 的背压策略
// BUFFER —— 默认,缓冲所有元素(可能导致 OOM)
Flux.just("a", "b", "c")
.onBackpressureBuffer(100)
.subscribe(subscriber);
// DROP —— 丢弃放不下的元素
Flux.interval(Duration.ofMillis(1))
.onBackpressureDrop(dropped ->
System.out.println("丢弃: " + dropped))
.subscribe(subscriber);
// LATEST —— 只保留最新元素,覆盖旧的
Flux.interval(Duration.ofMillis(1))
.onBackpressureLatest()
.subscribe(subscriber);
// ERROR —— 背压溢出时抛出异常
Flux.interval(Duration.ofMillis(1))
.onBackpressureError()
.subscribe(subscriber);5.4 limitRate —— 自动分批发请求
limitRate(n) 自动将上游请求拆分成 n 大小的小批量,降低上游缓冲压力:
Flux.range(1, 1000)
.limitRate(10) // 每次向下游请求 10 个,内部预取 75%
.subscribe(System.out::println);
// limitRate(10) 等价于:
// 首次 request(10),当已请求未消费数 <= 10×0.25=2.5 时再次 request(10)6. 错误处理与重试策略
6.1 静态回退
| 操作符 | 说明 |
|---|---|
onErrorReturn | 错误时返回一个固定值 |
onErrorResume | 错误时切换到备用 Publisher |
onErrorMap | 将异常转换为另一种异常后继续传播 |
// onErrorReturn
Flux.just("1", "2", "abc", "3")
.map(Integer::parseInt)
.onErrorReturn(-1)
.subscribe(System.out::println);
// 输出: 1, 2, -1
// onErrorMap —— 异常转换
Flux.just("ok")
.concatWith(Mono.error(new SQLException("连接超时")))
.onErrorMap(SQLException.class,
e -> new BusinessException("数据库异常", e))
.subscribe();6.2 重试策略(RetrySpec)
Reactor 3.3+ 提供了功能丰富的 RetrySpec:
// 基本重试:最多 3 次
Flux.just("data")
.concatWith(Mono.error(new RuntimeException("临时故障")))
.retryWhen(Retry.max(3))
.subscribe();
// 指数退避重试:初始 500ms,最大 10s,最多 5 次
Flux.just("data")
.concatWith(Mono.error(new RuntimeException("超时")))
.retryWhen(Retry.backoff(5, Duration.ofMillis(500))
.maxBackoff(Duration.ofSeconds(10))
.jitter(0.3)) // 抖动因子避免惊群
.subscribe();
// 按异常类型过滤重试
Flux.just("data")
.concatWith(Mono.error(new TimeoutException("超时")))
.retryWhen(Retry.max(3)
.filter(throwable -> throwable instanceof TimeoutException))
.subscribe();
// 重试间执行副作用(日志记录)
Flux.just("data")
.concatWith(Mono.error(new RuntimeException("错误")))
.retryWhen(Retry.max(2)
.doBeforeRetry(rs ->
System.out.println("第 " + (rs.totalRetries() + 1) + " 次重试")))
.subscribe();6.3 retryWhen 原理
retryWhen 的核心是一个 Flux<RetrySignal> 的伴随流。每次错误发生时,Reactor 向伴随流发射信号。伴随流上看到 onComplete 则重试,看到 onError 则停止。
// retryWhen 伴随流驱动
AtomicInteger retryCount = new AtomicInteger(0);
Flux.just("data")
.concatWith(Mono.error(new RuntimeException("fail")))
.retryWhen(Retry.from(companion ->
companion.take(3)
.delayElements(Duration.ofSeconds(1))
.doOnNext(s -> System.out.println("重试第 " + retryCount.incrementAndGet() + " 次"))
.concatWith(Mono.error(new RuntimeException("重试耗尽")))))
.subscribe();7. switchIfEmpty / defaultIfEmpty
这两个操作符用于处理空序列(上游没有发射任何元素就完成了)的场景。
7.1 defaultIfEmpty —— 提供默认值
// 场景:根据用户 ID 查询名称,找不到则返回默认
public Mono<String> getUserName(Long userId) {
return userRepository.findById(userId)
.map(User::getName)
.defaultIfEmpty("未知用户");
}
// 验证空序列
Mono<String> result = Mono.empty().defaultIfEmpty("默认值").block();
// result = "默认值"7.2 switchIfEmpty —— 切换到备用 Publisher
// 场景:先从本地缓存查询,缓存 miss 则查询数据库
public Mono<User> getUser(Long id) {
return cacheService.get(id)
.switchIfEmpty(Mono.defer(() ->
userDao.findById(id)
.flatMap(user ->
cacheService.set(id, user).thenReturn(user)
)
));
}
// 备用 Publisher 也可以是错误信号
Mono.empty()
.switchIfEmpty(Mono.error(new RuntimeException("无数据")))
.subscribe(); // 抛出异常7.3 defer 延迟求值
switchIfEmpty 内部应使用 Mono.defer(),避免不管是否为空都执行:
// 错误:switchIfEmpty 参数会立即执行
public Mono<User> badGetUser(Long id) {
return cacheService.get(id)
.switchIfEmpty(userDao.findById(id)); // 永远会被调用
}
// 正确:用 defer 延迟创建 Mono
public Mono<User> goodGetUser(Long id) {
return cacheService.get(id)
.switchIfEmpty(Mono.defer(() -> userDao.findById(id)));
}8. flatMap 并发控制(concurrency 参数)
8.1 默认行为
flatMap 默认的并发度是 256(Reactor 3.x),这意味着 flatMap 内部最多同时有 256 个内层 Publisher 处于订阅状态。某些场景下需要手动控制。
8.2 concurrency 参数的使用
// flatMap(Function, int concurrency, int prefetch)
Flux.range(1, 1000)
.flatMap(i -> asyncRemoteCall(i), 10, 32)
.subscribe();
// 实际场景:限制对外部 API 的并发调用
Flux.fromIterable(orderIds)
.flatMap(orderId ->
paymentService.refund(orderId)
.timeout(Duration.ofSeconds(5))
.onErrorReturn("退款失败: " + orderId),
5 // 同时最多 5 个退款请求并发
)
.subscribe(System.out::println);8.3 flatMapSequential —— 保持顺序的 flatMap
如果需要并发执行但保持原始顺序,使用 flatMapSequential:
Flux.range(1, 10)
.flatMapSequential(i ->
Mono.fromCallable(() -> remoteCall(i))
.subscribeOn(Schedulers.boundedElastic()),
3 // 并发度 3,但输出保持原始顺序
)
.subscribe(System.out::println);8.4 concatMap —— 严格有序(等效并发度为 1)
Flux.range(1, 5)
.concatMap(i -> Mono.delay(Duration.ofMillis(100))
.map(v -> "任务 " + i + " 完成"))
.subscribe(System.out::println);
// 输出严格有序,每个间隔 100ms9. 实战:高并发下 WebFlux + R2DBC 异步数据库连接池差异分析
9.1 背景
Spring WebFlux 配合 R2DBC(Reactive Relational Database Connectivity)可以实现从 HTTP 到数据库的全链路异步非阻塞。但在高并发场景下,连接池的行为与传统 JDBC 连接池存在显著差异。
9.2 连接池对比:HikariCP vs R2DBC Pool
| 维度 | HikariCP(JDBC) | R2DBC Pool(r2dbc-pool) |
|---|---|---|
| 线程模型 | 连接获取阻塞线程 | 非阻塞,返回 Mono<Connection> |
| 等待策略 | 调用线程阻塞等待 | 请求入队列,连接可用时异步回调 |
| 最大连接数 | 通常 10~50 | 通常 5~20 |
| 超时机制 | connectionTimeout(毫秒) | acquireTimeout(Duration) |
| 监控 | 自带 JMX | 需手动包装 Metrics |
9.3 高并发场景差异
场景:QPS 10000,每个请求执行 2 次数据库查询,连接池大小 10。
// === JDBC + HikariCP 的行为 ===
// 每个请求 -> 分配一个 Servlet 线程
// 线程 -> getConnection() -> 若池无可用连接则阻塞等待
// 后果:QPS 10000 需要大量线程,上下文切换开销巨大
// 配置:maximum-pool-size=50, connection-timeout=30000
// === R2DBC Pool 的行为 ===
// 每个请求 -> 不分配专用线程
// 请求连接 -> 返回 Mono<Connection>,调用线程不被阻塞
// -> 请求入队,event-loop 继续处理其他请求
// -> 连接可用时通过回调继续执行
// 后果:即使 QPS 10000,也只需少量线程处理回调
// 配置:max-size=15, max-acquire-time=5s9.4 关键风险与最佳实践
风险 1:事务中的连接泄漏
// 危险:事务中混用阻塞操作,连接被长时间占用导致池耗尽
@Transactional
public Mono<Void> processOrder(Order order) {
return orderRepo.save(order)
.flatMap(saved -> {
Thread.sleep(100); // 绝对禁止!
return Mono.empty();
});
}
// 正确:确保事务范围最小,避免阻塞
@Transactional
public Mono<Void> processOrder(Order order) {
return orderRepo.save(order).then();
}风险 2:连接池大小设置不当
// R2DBC 连接池大小经验公式:
// 连接数 ≈ 期望 QPS × 单次查询耗时
// 场景:QPS 5000,单次查询 5ms
// 需要连接数 ≈ 5000 × 0.005 = 25,加余量 25 × 1.5 ≈ 38
// 实际配置建议:
// spring.r2dbc.pool.max-size=30
// spring.r2dbc.pool.max-idle-time=10m
// spring.r2dbc.pool.max-life-time=30m风险 3:无超时保护的队列堆积
// acquireTimeout 设置过长会导致等待队列堆积,最终 OOM
spring:
r2dbc:
pool:
max-size: 20
max-acquire-time: 5s
max-create-connection-time: 5s
acquire-retry: 2
initial-size: 59.5 监控与调试
// 自定义 R2DBC 连接池度量(基于 Micrometer)
@Configuration
public class R2dbcPoolMetricsConfig {
@Bean
public ConnectionPoolMetrics connectionPoolMetrics(
@Qualifier("connectionFactory") ConnectionPool pool) {
return new ConnectionPoolMetrics(pool, "myapp.r2dbc.pool",
Tags.of("app", "order-service"));
}
static class ConnectionPoolMetrics {
ConnectionPoolMetrics(ConnectionPool pool, String prefix, Tags tags) {
MeterRegistry registry = Metrics.globalRegistry;
Gauge.builder(prefix + ".acquired", pool,
p -> p.getMetrics().getAcquiredSize())
.tags(tags).register(registry);
Gauge.builder(prefix + ".pending", pool,
p -> p.getMetrics().getPendingAcquireSize())
.tags(tags).register(registry);
Gauge.builder(prefix + ".idle", pool,
p -> p.getMetrics().getIdleSize())
.tags(tags).register(registry);
Gauge.builder(prefix + ".max", pool,
p -> p.getMetrics().getMaxAllocatedSize())
.tags(tags).register(registry);
}
}
}9.6 性能对比测试数据(参考)
| 连接池类型 | 线程数 | QPS | P99 延迟 | CPU 使用率 |
|---|---|---|---|---|
| HikariCP (30) | 200 (Tomcat) | 8500 | 45ms | 78% |
| R2DBC Pool (15) | 12 (Netty) | 9200 | 28ms | 52% |
| R2DBC Pool (30) | 12 (Netty) | 10500 | 22ms | 58% |
R2DBC Pool 在线程数极少的情况下实现了更高吞吐量,且 P99 延迟更低。这是全链路异步非阻塞带来的优势。
9.7 选型建议
- 低并发场景(QPS < 1000):JDBC + HikariCP 足够,开发维护成本更低
- 高并发 I/O 密集型(QPS > 5000):优先选择 WebFlux + R2DBC Pool
- 混合场景:WebFlux + JDBC 时,务必用
Schedulers.boundedElastic()隔离阻塞调用 - 连接池大小:R2DBC 连接池通常设得比 HikariCP 更小,因为线程等待开销极低
总结
本文从 Reactive Streams 规范出发,系统梳理了 Spring WebFlux 响应式编程的核心知识体系:
- Reactive Streams 定义了
Publisher-Subscriber-Subscription的背压契约 - Flux/Mono 操作符(map/flatMap/filter/zip/merge/concat)提供了声明式的数据流组合能力
- subscribe 的多种方式以及自定义
Subscriber实现对背压的精细控制 - Scheduler 调度器和
subscribeOn/publishOn决定了响应式链的线程模型 - 背压机制 是响应式流的灵魂,防止生产者压倒消费者
- 错误处理 通过
onErrorReturn/onErrorResume/retryWhen实现弹性容错 - switchIfEmpty / defaultIfEmpty 优雅处理空序列
- flatMap 并发控制 通过 concurrency 参数限制并发度
- R2DBC 连接池差异 分析了高并发下全链路异步的优势与风险
掌握这些基础,开发者就能在实际项目中合理运用 WebFlux 构建高性能、高弹性的响应式系统。