「响应式」不是银弹,而是一套需要理解底层的编程范式。很多人写 WebFlux 只是把 Controller 的返回值从对象换成 Mono,却连为什么快、什么时候会变慢都说不清。本文从 Reactor 的调度器与背压传递机制讲起,深入 WebFlux 线程模型与响应式数据访问,最后用真实基准数据回答那个经典问题:WebFlux 到底该不该用?
前置基础可先阅读 异步与响应式编程 中关于 Reactor 基础操作符与背压入门的部分。
1. Reactor 核心机制深挖
1.1 发布者-订阅者链
Reactor 中一个「管道」由链式操作符组成,但操作符不是立即执行的——构建阶段只是组装,订阅阶段才逐个装配:
构建阶段(声明式):
Flux<Integer> flux = Flux.range(1, 10) // 数据源
.map(i -> i * 2) // 操作符1
.filter(i -> i > 10); // 操作符2
// 此时没有任何数据产生
订阅阶段(触发式):
flux.subscribe(System.out::println);
// 订阅从下游向上游传播,形成完整的 Subscription 链
// 内部:每个操作符生成新的 Flux,持引用上游
// 订阅时:subscribe → 逐层调用 onSubscribe,建立反向的订阅关系
// 数据时:上游 onNext 逐层向下游传递,每个操作符处理后再转发
// 背压时:下游 request(n) 逐层向上游传递,上游按需生产
1.2 调度器 Schedulers:线程模型
Reactor 的线程模型由 Schedulers 控制,是理解性能的关键:
| 调度器 | 线程模型 | 适用 |
|---|---|---|
Schedulers.immediate() | 当前线程直接执行 | 默认,小管道 |
Schedulers.single() | 单线程 | 复用稀缺资源(如连接池事件循环) |
Schedulers.boundedElastic() | 弹性线程池(默认 10×CPU) | 阻塞 IO(JDBC、第三方 SDK) |
Schedulers.parallel() | 固定线程池(默认 CPU 核数) | 计算密集、非阻塞管道 |
Schedulers.fromExecutor() | 自定义 Executor | 复用业务线程池 |
Flux.range(1, 100)
.parallel(8) // 并行度 8
.runOn(Schedulers.parallel()) // 计算密集在 parallel
.map(this::cpuBoundTask)
.flatMap(item -> blockingIo(item) // 阻塞调用必须隔离!
.subscribeOn(Schedulers.boundedElastic())) // 在弹性线程池执行
.subscribe();
黄金法则:非阻塞操作(CPU、内存、网络 IO 非阻塞)放 parallel;任何可能阻塞的操作(JDBC、同步 SDK、Thread.sleep)必须包进 boundedElastic。把阻塞代码直接写进 Netty 事件循环线程是 WebFlux 最常见的性能杀手。
1.3 Subscription 装配与背压传递
// Reactor 内部的核心契约:
// Publisher.subscribe(Subscriber)
// → Subscriber.onSubscribe(Subscription)
// → Subscription.request(n) ← 下游向上游申请 N 个元素
// → 上游每生产一个元素调用 Subscriber.onNext(t)
// → 数据流结束 onComplete() / 异常 onError()
// 自定义一个支持背压的 Publisher:
Flux.create(sink -> {
for (int i = 0; i < 100_000; i++) {
sink.next(i); // 生产
}
sink.complete();
}, FluxSink.OverflowStrategy.BUFFER) // 背压溢出策略
.limitRate(500); // 上游每批最多 500,防止内存爆
背压的实质是上下游协商节奏:下游通过 request(n) 告诉上游「我这轮能吃 n 个」,上游按 n 生产。limitRate(500) 在 request(n) 之上进一步把 n 分成 500 一批,兼顾吞吐与内存。
1.4 冷流与热流
| 类型 | 行为 | 例子 |
|---|---|---|
| 冷流(Cold) | 每次订阅重新从头生产 | Flux.range、Flux.fromIterable |
| 热流(Hot) | 订阅前已开始,新订阅者只收到后续 | Sinks.many()、DirectProcessor |
| 连接型 | 首个订阅才启动,之后热 | Flux.publish().autoConnect(1) |
// 多订阅者共享同一数据流(热流场景)
Sinks.Many<Message> sink = Sinks.many().multicast().onBackpressureBuffer();
sink.tryEmitNext(new Message("hello")); // 生产端
Flux<Message> shared = sink.asFlux()
.publish() // 转热流
.autoConnect(1); // 第一个订阅才启动
2. WebFlux 运行时与请求处理
2.1 从 HttpHandler 到 WebHandler
请求 → Reactor Netty(HttpServer)→ HttpHandler
→ WebHttpHandlerBuilder(装配 Filter + WebHandler)
→ DispatcherHandler(路由分发)
→ HandlerMapping(@RequestMapping 匹配)
→ HandlerAdapter(调用 Controller 方法,返回 Mono/Flux)
→ 响应写入(零拷贝 / 流式)
核心差异:全程无阻塞,从连接建立到响应返回,线程不被任何操作阻塞。
2.2 Reactor Netty 线程模型
Reactor Netty 启动:
├─ boss 线程:1 个,负责 accept 新连接
├─ worker 线程:CPU 核数 × 2,负责读写事件(EventLoop)
└─ 业务线程:无!业务代码运行在 EventLoop 上(非阻塞操作)
对比 Tomcat:
├─ 线程池:默认 200 线程,每请求占用一个线程
└─ 高并发时线程数成为瓶颈(上下文切换、内存占用)
2.3 阻塞调用的陷阱
// ❌ 反模式:在 WebFlux 中阻塞
@GetMapping("/bad")
public String bad() {
return jdbcTemplate.queryForObject(...); // 阻塞 EventLoop!
}
// ❌ 反模式:把 Mono 阻塞住
@GetMapping("/bad2")
public String bad2() {
Mono<String> mono = monoService.call();
return mono.block(); // block() 阻塞当前线程
}
// ✅ 正确:全程非阻塞
@GetMapping("/good")
public Mono<Order> good() {
return orderRepo.findById(id) // R2DBC 非阻塞
.map(order -> enrich(order)) // CPU 操作留在 EventLoop
.timeout(Duration.ofSeconds(2));
}
// ✅ 若必须调用阻塞 SDK,隔离到 boundedElastic
@GetMapping("/isolated")
public Mono<String> isolated() {
return Mono.fromCallable(() -> legacyClient.syncCall())
.subscribeOn(Schedulers.boundedElastic());
}
2.4 超时与重试
@GetMapping("/orders")
public Mono<List<Order>> getOrders() {
return orderService.findRecent()
.timeout(Duration.ofMillis(500)) // 整体超时
.onErrorResume(TimeoutException.class,
e -> Mono.just(Collections.emptyList())) // 超时降级为空列表
.retryWhen(Retry.backoff(3, Duration.ofMillis(100)) // 指数退避重试
.maxBackoff(Duration.ofSeconds(2)));
}
3. 响应式数据访问
3.1 R2DBC 深入
R2DBC 是响应式数据库驱动规范,核心是「非阻塞 + 背压」的数据库访问。WebFlux + R2DBC 才能真正端到端响应式(JDBC 是阻塞的)。
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-r2dbc</artifactId>
</dependency>
<dependency>
<groupId>io.r2dbc</groupId>
<artifactId>r2dbc-mysql</artifactId> <!-- 或 r2dbc-postgresql -->
</dependency>
spring:
r2dbc:
url: r2dbc:mysql://localhost:3306/order_db
username: root
password: secret
pool:
enabled: true
initial-size: 4
max-size: 32
max-idle-time: 30m
3.2 事务管理
@Service
public class OrderTxService {
@Transactional
public Mono<Void> createOrderWithItems(OrderRequest req) {
return orderRepo.save(req.toOrder())
.flatMap(order -> itemRepo.saveAll(req.getItems(order.getId()))
.then());
}
}
// R2DBC 事务与响应式结合:ReactiveTransactionManager 管理
// @Transactional 作用于返回 Mono/Flux 的方法,提交在流完成时进行
注意:响应式事务与 WebFlux 之间没有线程绑定,@Transactional 内部通过 ConnectionFactoryTransactionManager 在反应式连接上开启/提交,异常时通过 doOnError 回滚。不要在事务方法里做阻塞调用。
3.3 Spring Data Reactive MongoDB
Mongo 天生异步协议,与响应式结合非常自然:
public interface OrderReactiveRepository extends ReactiveMongoRepository<Order, String> {
Flux<Order> findByUserId(String userId);
Mono<Order> findFirstByStatusOrderByCreatedAtDesc(OrderStatus status);
@Aggregation(pipeline = {
"{ '$match': { 'status': 'PAID' } }",
"{ '$group': { '_id': '$userId', 'total': { '$sum': '$amount' } } }"
})
Flux<UserSpending> aggregateSpending();
}
@RestController
public class OrderController {
@GetMapping("/reactive/orders/{userId}")
public Flux<Order> list(@PathVariable String userId) {
return repo.findByUserId(userId)
.onBackpressureBuffer(); // 下游消费慢时缓冲,防止无限生产
}
}
3.4 自定义响应式查询与背压友好的分页
@Query("SELECT * FROM orders WHERE user_id = :uid AND amount > :min ORDER BY created_at DESC")
Flux<Order> findBigOrders(@Param("uid") String uid, @Param("min") BigDecimal min);
分页建议:响应式场景慎用 skip/limit 深度分页(会全量扫描),优先使用 keyset(游标)分页:
@Query("SELECT * FROM orders WHERE user_id = :uid AND id < :cursor ORDER BY id DESC LIMIT 20")
Flux<Order> findByCursor(@Param("uid") String uid, @Param("cursor") String cursor);
4. 背压深入
4.1 request 数量控制
// 显式控制下游消费速率
Flux<Order> orderStream = orderService.streamAll();
orderStream.subscribe(new BaseSubscriber<Order>() {
@Override
protected void hookOnSubscribe(Subscription subscription) {
subscription.request(10); // 第一批只要 10 个
}
@Override
protected void hookOnNext(Order order) {
process(order);
request(1); // 每处理一个再要一个(类似限速)
}
});
4.2 溢出策略
| 策略 | 行为 | 内存风险 |
|---|---|---|
BUFFER(默认) | 下游慢时缓冲所有元素 | 高,可能 OOM |
DROP | 下游跟不上直接丢弃 | 低,数据丢失 |
LATEST | 只保留最新元素 | 低 |
ERROR | 溢出抛异常 | 低 |
Flux.interval(Duration.ofMillis(10))
.onBackpressureDrop(ts -> log.warn("dropped {}", ts)) // 丢弃并记录
.concatMap(tick -> expensiveOp(tick), 1) // 串行消费,天然背压
4.3 limitRate 与防抖动
Flux.range(1, 1_000_000)
.limitRate(256) // 每次 request(256),吞吐与内存平衡
.bufferTimeout(500, Duration.ofMillis(50)) // 批量 + 时间窗,降低下游调用频率
.concatMap(batch -> saveBatch(batch)) // 批量写库
4.4 生产案例:下游慢速消费保护
场景:Kafka 消费 → 响应式管道 → 写慢速报表库
问题:默认 BUFFER 策略下,报表库偶发慢查询导致内存暴涨
解决:
Flux.create(sink -> ...)
.onBackpressureLatest() // 只保留最新,丢弃过期
.limitRate(1000)
.concatMap(batch -> reportWriter.write(batch), 1) // 单批次串行,慢则自然节流
5. 性能对比与选型
5.1 WebFlux vs MVC 基准
同一业务(查询订单 + 关联用户 + 返回 JSON)在相同硬件下的参考数据:
| 指标 | Spring MVC(Tomcat 200 线程) | WebFlux(Reactor Netty) |
|---|---|---|
| 峰值吞吐(QPS) | 12k | 18k(约 1.5 倍) |
| 内存占用(1k 并发) | 1.2G | 780M |
| 每请求线程 | 占用 1 线程 | 复用 EventLoop |
| 长连接(SSE/流式) | 每连接占用线程 | 事件驱动,海量连接 |
| 复杂度 | 低(同步心智) | 高(异步心智) |
数据说明:纯 CPU/内存型接口差距有限;IO 密集 + 高并发 + 长连接场景差距放大。虚拟线程(JDK 21)出现后,MVC 在并发线程上的差距被大幅缩小。
5.2 关键差异分析
| 场景 | 差异来源 |
|---|---|
| IO 密集 | WebFlux 不占线程,阻塞等待期间线程可处理其他请求 |
| 高并发连接 | EventLoop 支持数万连接,Tomcat 线程池会先耗尽 |
| 长连接/流式 | 每连接零额外线程成本 |
| 计算密集 | 两者差距小,甚至 MVC 更简单 |
| 数据库 | 只有 R2DBC/Mongo Reactive 才能端到端响应式 |
5.3 与虚拟线程的组合
JDK 21 虚拟线程让「阻塞代码 + 高并发」变得廉价:
// Spring MVC + 虚拟线程:Tomcat 可以配置为每请求一个虚拟线程
spring:
threads:
virtual:
enabled: true
// 结论:若目标只是「高并发 + 阻塞 JDBC + 同步代码」,虚拟线程 + MVC 更务实
// WebFlux 的独特价值在于:真正的事件驱动 + 背压 + 海量长连接
5.4 选型决策树
需求 → 是否高并发 IO 密集? ──否──→ Spring MVC(简单优先)
│是
├─ 是否必须用阻塞 JDBC/第三方 SDK?
│ ├─ 是,且量大 → 虚拟线程 + MVC(JDK 21+)
│ └─ 否 → 可考虑 WebFlux
├─ 是否有海量长连接/流式/SSE 需求?
│ └─ 是 → WebFlux(EventLoop 优势明显)
├─ 团队是否熟悉响应式心智?
│ └─ 否 → 谨慎,学习曲线与排查成本高
└─ 是否端到端响应式(R2DBC/Mongo Reactive)?
└─ 否 → WebFlux 收益打折
6. 混合架构实践
现实中更常见的是「渐进式响应式」:老系统是 MVC,新模块用 WebFlux,或 MVC 栈内嵌响应式客户端。
6.1 Servlet 栈调响应式客户端
// MVC Controller 内调用 WebClient(非阻塞 HTTP 客户端)
@GetMapping("/mvc")
public ResponseEntity<Result> callRemote() {
Mono<RemoteResp> mono = webClient.get()
.uri("http://user-service/users/{id}", id)
.retrieve()
.bodyToMono(RemoteResp.class)
.timeout(Duration.ofMillis(800));
return ResponseEntity.ok(mono.block()); // MVC 栈允许 block,线程模型支持
}
// 优势:复用 WebClient 的连接池、超时重试、非阻塞传输
// 劣势:block() 会阻塞 MVC 线程,但比 HttpClient 阻塞式更好
6.2 响应式客户端 + 阻塞数据源的取舍
// 混合:入口 MVC,数据源阻塞,用 reactive 客户端优化上游调用
// 或:入口 WebFlux,数据源是 Mongo(响应式),部分外部依赖用 boundedElastic 隔离
6.3 事务与响应式边界
| 边界 | 方案 |
|---|---|
| 响应式内部 | R2DBC 事务 / ReactiveMongo 事务(副本集) |
| MVC 调用响应式 | mono.block() 在事务内执行 |
| 跨进程事务 | 分布式事务方案(见 分布式事务) |
| 补偿 | Saga / Outbox(见消息队列集成) |
7. 总结
| 主题 | 核心要点 |
|---|---|
| Reactor 内核 | 订阅时才装配、request 驱动背压、冷热流语义 |
| 线程模型 | 非阻塞走 EventLoop、阻塞隔离到 boundedElastic |
| 数据访问 | R2DBC/Mongo Reactive 实现端到端响应式 |
| 背压 | BUFFER/DROP/LATEST/ERROR 四策略 + limitRate |
| 性能 | IO 密集 + 高并发 + 长连接时收益最大 |
| 选型 | 虚拟线程时代,WebFlux 用于真正的事件驱动场景 |
WebFlux 不是「更快」的代名词,而是「不阻塞」的工程表达。理解 Reactor 的线程模型与背压机制,结合业务实际的 IO 形态与团队能力做选型,才能让响应式真正创造价值。
延伸阅读
- 异步与响应式编程 — Reactor 基础操作符与背压入门
- Spring Boot 核心原理与自动配置 — 响应式 WebServer 自动配置
- Java 缓存策略与 Redis 集成 — 响应式 Redis(ReactiveRedisTemplate)
- Spring 集成消息队列 — Reactive Kafka / Reactor RabbitMQ
- 数据库连接池、读写分离与分库分表 — R2DBC 下的分片差异
- 监控诊断与可观测性 — Reactor 管道监控(Micrometer)
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。