Administrator
发布于 2021-09-18 / 1390 阅读
37

WebFlux 响应式编程:什么场景下才值得用

背景:一个新服务要不要上 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 时 QPSP99峰值线程数峰值堆
MVC + CompletableFuture(Tomcat 200 线程)3,240382 ms4121.8 GB
WebFlux + WebClient9,76096 ms410.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。

方案QPSP99峰值线程数
MVC + 阻塞 JDBC2,910428 ms408
WebFlux + JDBC 卡在 EventLoop6206,410 ms10
WebFlux + JDBC 切 boundedElastic3,060394 ms68

关键数据在这:切到 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 响应式编程:什么场景下才值得用》里同一个坑绊住,回来翻这篇,能省半小时。

参考