背景:一个新服务要不要上 WebFlux
9 月中旬要起一个新服务,作用是"详情页聚合":前端调一次接口,我们并行调商品、库存、价格、营销、评价五个下游,聚合后返回。组里有人提议用 WebFlux,理由是"调用多、IO 密集,响应式最合适"。
我当时持保留态度,就花了一周做了个 POC,两边都实现一遍再压测。结论是:这个服务不适合上 WebFlux,而原因不是"响应式不好",是 JDBC。
先把线程模型说清楚
Spring MVC(Servlet 栈)的模型是"一个请求占一个线程":Tomcat 默认 200 个工作线程,请求进来从线程池取一个线程,控制器里的所有代码(查库、调 HTTP、处理结果)都在这个线程上跑,直到返回才归还。如果一个请求要等 5 个下游各 50 ms,这个线程就被占 50 ms 以上。
并发 400 的时候,200 个线程全在等,另外 200 个请求在队列里排队。这就是为什么阻塞模型在 IO 密集场景下吞吐上不去——线程数就是并发上限。
WebFlux 的模型是事件循环:底层是 Netty 的 EventLoop(默认线程数 2 * CPU 核数),所有 IO 操作都注册为非阻塞,等数据就绪时由事件循环回调处理。一个线程可以同时"持有"成百上千个等待中的请求,因为等待不占用线程,只是内存里保留了一个回调对象。
这就是吞吐差异的来源。理论很好,但有个前提:整条链路上不能有任何阻塞调用。
背压:消费者说了算
Reactive Streams 规范定义了四个接口:Publisher(生产者)、Subscriber(消费者)、Subscription(订阅关系)、Processor(既是生产者也是消费者)。背压的核心在 Subscription:
public interface Subscription {
void request(long n); // 消费者告诉生产者:我还能处理 n 个
void cancel();
}
public interface Subscriber<T> {
void onSubscribe(Subscription s);
void onNext(T t);
void onError(Throwable t);
void onComplete();
}
request(n) 是拉取式的:不是生产者推多少消费者就得接多少,而是消费者主动声明"我还要 n 个",生产者最多发 n 个。这样慢消费者不会被快生产者压垮,也就不需要无界队列。
在 WebFlux 里这个过程是自动的。举个例子,读一个 2 GB 的文件返回给客户端:
@GetMapping(value = "/download", produces = MediaType.APPLICATION_OCTET_STREAM_VALUE)
public Flux<DataBuffer> download() {
// 客户端慢,Netty 就不会从磁盘继续读,内存占用恒定
return DataBufferUtils.read(inputStreamSupplier, dataBufferFactory, 8192);
}
如果客户端网络慢,request(n) 的 n 就会变小,文件读取自动减速。同样的需求在 Servlet 栈里要么全读进内存(OOM),要么手写 OutputStream 分块写并处理背压。
中间的操作符也会影响背压传递。flatMap(f, concurrency) 里的 concurrency 就是一次向上游 request 多少个(默认 256);onBackpressureBuffer()、onBackpressureDrop()、onBackpressureLatest() 是三种处理不了时的策略,默认是无界缓冲,也就是会 OOM。
压测:两个场景,两个结论
POC 我写了两份实现,压测环境 4 核 8 G,JMeter 压 5 分钟。
场景一:纯 IO 聚合(下游用 Mock,固定延迟 50 ms)
// WebFlux 版:并行调用 5 个下游
@GetMapping("/detail/{skuId}")
public Mono<SkuDetailVO> detail(@PathVariable Long skuId) {
Mono<ItemInfo> item = itemClient.getItem(skuId);
Mono<Stock> stock = stockClient.getStock(skuId);
Mono<Price> price = priceClient.getPrice(skuId);
Mono<Promo> promo = promoClient.getPromo(skuId);
Mono<Review> review = reviewClient.getReview(skuId);
// zip 会并行发起,全部完成后聚合
return Mono.zip(item, stock, price, promo, review)
.map(t -> SkuDetailVO.of(t.getT1(), t.getT2(), t.getT3(),
t.getT4(), t.getT5()))
.timeout(Duration.ofMillis(800));
}
// MVC 版:用 CompletableFuture 并行
@GetMapping("/detail/{skuId}")
public SkuDetailVO detail(@PathVariable Long skuId) {
CompletableFuture<ItemInfo> item = itemClient.getItemAsync(skuId);
CompletableFuture<Stock> stock = stockClient.getStockAsync(skuId);
// ... 五个
CompletableFuture.allOf(item, stock, price, promo, review).join();
return SkuDetailVO.of(item.join(), stock.join(), price.join(),
promo.join(), review.join());
}
| 方案 | 并发 500 时 QPS | P99 | 峰值线程数 | 峰值堆 |
|---|---|---|---|---|
| MVC + CompletableFuture(Tomcat 200 线程) | 3,240 | 382 ms | 412 | 1.8 GB |
| WebFlux + WebClient | 9,760 | 96 ms | 41 | 0.9 GB |
WebFlux 吞吐是 3 倍,线程数从 412 降到 41,堆占用减半。这个场景响应式完胜,符合预期。
场景二:加上数据库查询
但真实的服务不可能不查库。我们的聚合逻辑里有一条:从 t_sku_ext 表读商品的扩展属性,MySQL 8.0,单次查询 18 ms。加上之后:
// 错误写法:直接在响应式链上用阻塞的 JDBC
@GetMapping("/detail/{skuId}")
public Mono<SkuDetailVO> detail(@PathVariable Long skuId) {
return Mono.zip(item, stock, price, promo, review)
.map(t -> {
SkuExt ext = skuExtMapper.selectById(skuId); // ← 阻塞!
return SkuDetailVO.of(..., ext);
});
}
这段能跑,但 skuExtMapper.selectById 会阻塞 Netty 的 EventLoop 线程。4 核机器只有 8 个 EventLoop 线程,一个查库 18 ms 就把整条事件循环卡住 18 ms,所有请求都停了。压测结果是 QPS 掉到 620,P99 涨到 6.4 秒,比 MVC 差了一个数量级。
正确写法是把阻塞调用切到专门的线程池:
private static final Scheduler JDBC_SCHEDULER =
Schedulers.newBoundedElastic(50, 500, "jdbc-pool", 60);
return Mono.zip(item, stock, price, promo, review)
.publishOn(JDBC_SCHEDULER) // 切换到阻塞池
.map(t -> {
SkuExt ext = skuExtMapper.selectById(skuId);
return SkuDetailVO.of(..., ext);
})
.publishOn(Schedulers.parallel()); // 切回非阻塞池
Schedulers.boundedElastic() 是 Spring 提供的"为阻塞任务准备的"调度器:线程池按需增长,有上限(默认 10 * CPU 核数),空闲 60 秒回收。切过去之后阻塞调用不再占 EventLoop。
| 方案 | QPS | P99 | 峰值线程数 |
|---|---|---|---|
| MVC + 阻塞 JDBC | 2,910 | 428 ms | 408 |
| WebFlux + JDBC 卡在 EventLoop | 620 | 6,410 ms | 10 |
| WebFlux + JDBC 切 boundedElastic | 3,060 | 394 ms | 68 |
关键数据在这:切到 boundedElastic 之后,WebFlux 只比 MVC 快 5%。因为瓶颈已经变成那 50 个 JDBC 线程了——只要你用了阻塞的 JDBC,不管外面套什么模型,并发上限还是线程池大小。响应式的优势荡然无存,还多了一堆心智负担。
JDBC 是绕不开的现实约束
核心问题:JDBC 规范本身是阻塞的。Connection.prepareStatement()、ResultSet.next() 全是同步方法,没有标准的异步版本。JDBC 工作组提过 ADBC(Asynchronous Database Connectivity),但一直没落地。
替代方案是 R2DBC(Reactive Relational Database Connectivity),1.0.0 在 2021 年 5 月发布。我试过 r2dbc-mysql 0.8.2 + r2dbc-pool 0.9,能跑,纯响应式查库确实能回到 9800 QPS 的水平。但坑不少:
- 没有成熟的 ORM。Spring Data R2DBC 提供了
DatabaseClient和 Repository 支持,但比 MyBatis 差太远,复杂查询要自己拼Criteria或者裸 SQL。 - 不支持延迟加载、不支持一级缓存这类 MyBatis 特性。
- 事务要靠
TransactionalOperator或者@Transactional+ Reactor 上下文传播,排查问题时比DataSourceTransactionManager难懂。 - MySQL 驱动实现当时还是 0.8.x,没到 GA,我们不敢用在生产。
如果整个服务都是"转发 + 聚合",不怎么查库,那 WebFlux 很香。但我们的业务系统 CRUD 占大头,这个约束是硬的。
其他几笔隐性成本
除了 JDBC,还有几件事压测数据里看不出来,但实打实影响效率:
调试和栈信息。响应式链的异常栈长这样:
java.lang.NullPointerException: The mapper [xxx] returned a null value.
at reactor.core.publisher.FluxMapFuseable$MapFuseableSubscriber.onNext(FluxMapFuseable.java:113)
Suppressed: reactor.core.publisher.FluxOnAssembly$OnAssemblyException:
Assembly trace from producer [reactor.core.publisher.MonoFlatMap] :
reactor.core.publisher.Mono.flatMap(Mono.java:3089)
com.xxx.aggregate.SkuAggregateService.lambda$detail$3(SkuAggregateService.java:71)
Error has been observed at the following site(s):
|_ Mono.flatMap ⇢ at com.xxx.aggregate.SkuAggregateService.lambda$detail$3(SkuAggregateService.java:71)
|_ Mono.zip ⇢ at com.xxx.aggregate.SkuAggregateService.detail(SkuAggregateService.java:64)
看着信息多,但和"哪一行代码出错了"是两回事。要定位得靠 .checkpoint("读库存") 手工加检查点,或者开 Hooks.onOperatorDebug()(有性能损耗,只能调试时开)。团队里不熟悉的人上手至少要两周。
ThreadLocal 全部失效。MDC 日志链路、SkyWalking 的 traceId、用户上下文,在响应式链上都不能用 ThreadLocal 存,要改成 Reactor 的 Context:
// 写入
return chain.filter(exchange)
.contextWrite(ctx -> ctx.put("traceId", traceId));
// 读取(在 Mono/Flux 的操作符里)
Mono.deferContextual(ctx -> {
String traceId = ctx.get("traceId");
MDC.put("traceId", traceId);
return doSomething();
});
我们有一批公共库(权限校验、操作日志切面)深度依赖 ThreadLocal,全迁响应式意味着全部重写。这笔账最后是我们放弃的主因之一。
生态适配。不是所有中间件都有响应式驱动。Lettuce(Redis)有、WebClient(HTTP)有、MongoDB 有、Cassandra 有,但 Kafka、Elasticsearch、Dubbo 的响应式支持都不完整,用的时候还是得 publishOn 切线程。
一个隐蔽的雷:在响应式代码里调 .block() 会直接抛异常:
java.lang.IllegalStateException: block()/blockFirst()/blockLast() are blocking,
which is not supported in thread reactor-http-nio-2
这个设计是好的(阻止你把响应式代码写回阻塞),但新手经常中招,尤其是从老代码搬逻辑过来的时候。
WebClient 有几个参数必须配
就算不上 WebFlux,WebClient 作为 HTTP 客户端也值得用(Spring 5 已经把 RestTemplate 标记为"未来不再增加新特性")。但我们第一次用的时候踩了连接池的坑:
WebClient.builder()
.clientConnector(new ReactorClientHttpConnector(
HttpClient.create(ConnectionProvider.builder("custom")
.maxConnections(500) // 默认只有 2 * CPU
.maxIdleTime(Duration.ofSeconds(60))
.maxLifeTime(Duration.ofMinutes(10))
.pendingAcquireTimeout(Duration.ofSeconds(10)) // 默认 45 秒,太长
.evictInBackground(Duration.ofSeconds(120))
.build())))
.build();
默认的 ConnectionProvider 最大连接数是 2 * CPU 核数,4 核机器上就是 8 个。我们上线之后并发一上来,请求全卡在等连接上,pendingAcquireTimeout 又默认 45 秒,直接把 P99 拉到 40 秒。默认值是给"客户端只用几个连接"的场景设计的,服务端之间调用必须自己配。
另外超时一定要显式设,否则响应体读取会用到默认的 responseTimeout(无超时):
HttpClient.create()
.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 2000)
.responseTimeout(Duration.ofSeconds(3))
.doOnConnected(conn -> conn
.addHandlerLast(new ReadTimeoutHandler(3, TimeUnit.SECONDS))
.addHandlerLast(new WriteTimeoutHandler(3, TimeUnit.SECONDS)));
那什么场景该用
我的判断标准其实就一条:服务是不是"几乎不碰数据库",且并发高、下游慢。
| 适合 | 不适合 |
|---|---|
| API 网关(Spring Cloud Gateway 本身就是 WebFlux) | 以 CRUD 为主的业务服务 |
| 聚合层 / BFF,调用多个慢下游 | 重度依赖 JDBC 和 ORM 的系统 |
| SSE、WebSocket、长连接推送 | 团队对响应式不熟、没人能兜底 |
| 大文件流式传输、流式响应 | 深度依赖 ThreadLocal / MDC 的老代码 |
| 需要背压保护的下游(防止打爆) | 要求快速交付、工期紧的项目 |
我们最后这个聚合服务用了 Spring MVC + CompletableFuture 并行调用,Tomcat 线程调到 400,压测 4,100 QPS、P99 210 ms。比 WebFlux 版差一半,但开发效率、可维护性、排查难度上都划算得多。而且 4100 QPS 对这个服务的量级(峰值 800)已经严重过剩。
就写到这。如果哪天你也被《WebFlux 响应式编程:什么场景下才值得用》里同一个坑绊住,回来翻这篇,能省半小时。