WebFlux 响应式编程实战:Reactor 内核、响应式数据访问与选型

深入 Reactor 调度器与背压传递机制,掌握 WebFlux 线程模型、R2DBC/Mongo Reactive 数据访问,用基准数据对比 WebFlux 与 MVC 并给出选型决策树

「响应式」不是银弹,而是一套需要理解底层的编程范式。很多人写 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)12k18k(约 1.5 倍)
内存占用(1k 并发)1.2G780M
每请求线程占用 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 形态与团队能力做选型,才能让响应式真正创造价值。

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java-enterprise」更多文章

  1. 云原生 Java:GraalVM 原生镜像、镜像瘦身与 Serverless
  2. Java 安全与合规:安全编码、数据脱敏与供应链防护
  3. Java 微服务治理深化:熔断限流降级、灰度发布与全链路压测