楼层: 首页/ 软件技术/ JavaEE / Jakarta EE/ WebFlux 与响应式编程:高并发下的另一条路
11

WebFlux 与响应式编程:高并发下的另一条路

Reactive Stack · Reactor / Mono & Flux / WebClient / R2DBC

前面十章讲的都是 Servlet 那一套(阻塞式):一个请求占一个线程,线程等 IO 的时候就在那儿干等。这套模型简单直观,绝大多数业务够用。但当并发从几百涨到几万,问题就来了——线程是稀缺资源,一个线程占几 MB 栈内存,等 IO 的线程全是浪费。WebFlux 是 Spring 给出的另一条路:用少量线程处理海量请求,用"回调链"代替"阻塞等待"。这一章把它的执行模型、Reactor 核心 API、以及"什么时候千万别用"讲清楚。

两条路线的执行模型:线程到底在干什么

理解 WebFlux 的关键只有一句话:Servlet 是"一个请求一个线程,线程阻塞等 IO",WebFlux 是"少量线程轮转处理,IO 完成后再回来接着处理"。下面用一张表把这个差别量化。

维度Spring MVC(Servlet,阻塞)Spring WebFlux(Reactive,非阻塞)
IO 模型BIO。线程发起 IO 后就阻塞,直到数据返回NIO。注册回调,IO 就绪时事件循环再调度处理
线程模型Tomcat 线程池,默认 200 个线程。一个请求占一个线程直到响应完成Netty 事件循环,线程数约等于 CPU 核数(常见 4~16 个)
200 并发下的线程数约 200 个线程。若栈按 1MB 算,仅栈内存就约 200MB约 8~16 个线程,栈内存约 16MB,其余都在处理业务
并发上限由什么决定线程池大小。超出就排队,队列满了就拒绝内存与下游承载能力。几乎不因自身线程数而受限
下游慢时的表现线程被占住,池子很快耗尽,新请求全部排队(雪崩起点)线程立即释放去做别的事,等下游响应回来再继续
调试难度低。断点、调用栈都直观,异常栈清晰高。调用栈里全是 Reactor 的框架帧,业务逻辑散在 lambda 里
学习成本低。会写 Java 就会写高。要理解 Publisher/Subscriber、背压、调度器

论为什么"省线程"这件事这么重要

① 线程不是免费的:每个线程有自己的栈(默认 1MB,容器里常见 512KB),还有创建、销毁、上下文切换的开销。几千个线程的光是上下文切换就能吃掉大量 CPU。

② 阻塞式模型的真正问题不是"慢",而是"浪费":一个请求调下游接口要 200ms,这 200ms 里那个线程什么都没做,只是举着电话等对方说话。如果 QPS 是 5000,你就需要 1000 个线程同时在线——而其中 99% 的时间它们在等待。

③ WebFlux 的账是这样算的:同样 5000 QPS,用 8 个事件循环线程 + 异步 IO,就能全部处理完。线程数从 1000 降到 8,省的是内存、上下文切换和调度成本。前提是:链路里不能有任何阻塞操作。

Reactor 核心:Mono、Flux 与那几个最容易搞混的操作符

WebFlux 默认使用 Reactor 作为响应式库。它的两个核心类型很朴素:Mono<T> 表示"0 或 1 个结果",Flux<T> 表示"0 到 N 个结果"。它们都是"发布者(Publisher)",在你订阅之前什么都不做。

先跑一遍最基础的:创建、变换、消费

import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

public class ReactorBasics {

    public static void main(String[] args) {
        // ---------- Mono:0 或 1 个元素 ----------
        Mono<String> mono = Mono.just("hello");
        Mono<String> empty = Mono.empty();
        Mono<String> error = Mono.error(new RuntimeException("炸了"));

        // 关键:不订阅,什么都不会发生!下面这行不会打印任何东西
        mono.map(String::toUpperCase);

        // 订阅之后才真正执行
        mono.map(String::toUpperCase)
            .subscribe(
                value -> System.out.println("收到: " + value),
                err   -> System.out.println("出错: " + err.getMessage()),
                ()    -> System.out.println("完成")
            );

        // ---------- Flux:0 到 N 个元素 ----------
        Flux<Integer> nums = Flux.just(1, 2, 3, 4, 5);

        nums.filter(n -> n % 2 == 0)      // 只要偶数
            .map(n -> n * 10)             // 每个都乘 10
            .subscribe(System.out::println);  // 20 40

        // 定时产生数据:每秒一个,共 5 个(常用于测试或心跳场景)
        Flux.interval(java.time.Duration.ofSeconds(1))
            .take(5)
            .subscribe(i -> System.out.println("tick " + i));
    }
}
操作符语义与选择依据
map同步变换:把 A 变成 B,一步完成,不产生新的 Publisher。返回普通对象时用这个。
flatMap异步变换:把 A 变成 Publisher<B>,并发合并结果,顺序不保证。适合"每个元素都要调一次远程接口、且不在意顺序"。
concatMap异步但串行:同样把 A 变成 Publisher<B>,但严格按顺序、一个一个来。适合"必须保序"的场景,代价是吞吐低于 flatMap。
zip把多个 Publisher 的第 N 个元素配对组合(像一个拉链)。适合"同时查用户和订单,两个都回来了再拼结果"。
merge把多个 Publisher 交织在一起输出,不等待配对,谁先来谁先出。适合"同时发两个请求,谁先返回就先处理"。
zipWith / then / thenReturn两个流配对 / 忽略前一个流的值只等它完成 / 等完成后返回固定值。常用于"先保存再返回固定响应"。

flatMap、concatMap、zip 的差别,用一段代码看清

// 假设有个远程接口:根据 id 查用户名,耗时 100ms
private Mono<String> queryNameAsync(Long id) {
    return Mono.fromCallable(() -> {
        Thread.sleep(100);              // 模拟远程调用(真实场景用 WebClient)
        return "user-" + id;
    });
}

// ---------- flatMap:并发执行,总耗时约 100ms,但顺序不保证 ----------
Flux.just(1L, 2L, 3L)
    .flatMap(id -> queryNameAsync(id))
    .subscribe(System.out::println);     // 可能是 user-2 user-1 user-3

// ---------- concatMap:串行执行,总耗时约 300ms,顺序严格保证 ----------
Flux.just(1L, 2L, 3L)
    .concatMap(id -> queryNameAsync(id))
    .subscribe(System.out::println);     // 一定按 user-1 user-2 user-3 顺序,但更慢

// ---------- zip:两个流配对,等到两边都有第一个元素才开始输出 ----------
Mono<String> user  = queryNameAsync(1L);
Mono<Integer> score = Mono.just(90);

Mono<String> combined = Mono.zip(user, score)
    .map(t -> t.getT1() + " 得分 " + t.getT2());
combined.subscribe(System.out::println);   // user-1 得分 90

// ---------- 注意:zip 有"短板效应" ----------
// 如果一边只发 3 个元素、另一边只发 2 个,zip 结果只有 2 个
// 如果一边永远不发(比如 Mono.never()),zip 会一直等着,表现为"请求卡住"

论背压(Backpressure):生产者别把消费者淹死

① 问题场景:上游(比如读一个大文件、或者一个高速消息队列)产生数据的速度远快于下游处理速度。如果没有机制,数据会在内存里堆积,最终 OOM。

② 背压的本质:让消费者告诉生产者"我还能接多少"。Reactor 的 subscribe 默认请求 Long.MAX_VALUE(也就是"我全要"),但你可以用 BaseSubscriber 精确控制请求数量,实现"处理完一个再要一个"。

③ 现实中的处理策略:不是所有场景都要严格的背压。常见三种做法:限流(limitRate,一次只请求 N 个)、丢弃(onBackpressureDrop,接不住就扔,适合实时性要求高的指标上报)、缓冲(onBackpressureBuffer,先存着,但要设上限防止 OOM)。

④ 提醒一句:WebFlux 对 HTTP 请求体是有背压的(TCP 天然有),但如果你在链里跳过了背压(比如用了 Flux.create 却不管请求量),照样会出问题。

subscribeOn 与 publishOn:调度到底在哪一步发生

这是响应式编程最容易搞混的一对。用一句话记住:subscribeOn 决定"从哪开始执行"(影响订阅信号向上传播),publishOn 决定"从哪之后在哪个线程执行"(影响数据向下游传播)。

用打印出来的线程名把语义看清楚

import reactor.core.scheduler.Schedulers;

// ---------- 场景一:subscribeOn 决定"源头在哪个线程执行" ----------
Mono.fromCallable(() -> {
        System.out.println("执行线程: " + Thread.currentThread().getName());
        return "data";
    })
    .subscribeOn(Schedulers.boundedElastic())     // 让订阅动作在弹性线程池执行
    .subscribe();
// 输出:执行线程: boundedElastic-1
// subscribeOn 的位置不重要,它影响的是"整个链从哪开始跑"

// ---------- 场景二:publishOn 切换"下游在哪执行" ----------
Mono.just("a")
    .map(v -> { System.out.println("map1: " + Thread.currentThread().getName()); return v; })
    .publishOn(Schedulers.parallel())              // 从这行之后换线程
    .map(v -> { System.out.println("map2: " + Thread.currentThread().getName()); return v; })
    .subscribe();
// 输出:
//   map1: main                  ← publishOn 之前,还是 main 线程
//   map2: parallel-1            ← publishOn 之后,切到了 parallel 线程

// ---------- 调度器的选择 ----------
// immediate()      :当前线程,不切换。适合"已经是异步的了,不用再切"
// single()         :单个可复用线程。适合"必须串行、但有顺序要求"
// parallel()       :固定大小线程池(约等于 CPU 核数)。适合纯 CPU 计算
// boundedElastic() :弹性线程池,有上限。适合阻塞式 IO(老 JDBC、文件 IO)
// ==========> 记住:CPU 密集用 parallel,阻塞调用用 boundedElastic

// ---------- 最标准的写法:把阻塞调用隔离到 boundedElastic ----------
public Mono<String> callLegacyBlockingApi(String param) {
    return Mono.fromCallable(() -> legacyService.blockingCall(param))
               .subscribeOn(Schedulers.boundedElastic());
    // 这样阻塞调用被挪到了弹性线程池,不会卡住 Netty 的事件循环
    // 这是"必须调用阻塞 API 时"唯一的正确姿势
}
subscribeOn / publishOn 的三个误区

误区一:以为写多个 subscribeOn 能多次切换线程。不能。整条链上只有第一个 subscribeOn 生效,后面的会被忽略。要多次切换必须用 publishOn。

误区二:以为 publishOn 影响上游。不能。publishOn 只影响它下游的操作符,上游该在哪个线程还是在哪个线程。所以"上游是阻塞调用"这种情况,用 publishOn 是治不了的——必须用 subscribeOn。

误区三:随便用 parallel() 跑阻塞 IO。parallel 的线程数只有 CPU 核数(比如 8 个)。如果你在里面调用会阻塞 200ms 的 JDBC,8 个线程几毫秒就被占满,整个应用就卡死了——比 MVC 还惨,因为 MVC 至少有 200 个线程。阻塞调用一律用 boundedElastic()。

WebClient 与函数式端点

进入实践部分。WebFlux 生态里,所有"调用别人的 HTTP 接口"都应该用 WebClient 而不是 RestTemplate,因为 RestTemplate 是阻塞的——在事件循环线程里调用它,等于把整条事件循环卡住。

WebClient 替代 RestTemplate:一次配置,处处异步调用

import org.springframework.web.reactive.function.client.WebClient;
import reactor.core.publisher.Mono;
import java.time.Duration;

@Configuration
public class WebClientConfig {

    @Bean
    public WebClient userWebClient() {
        return WebClient.builder()
            .baseUrl("http://user-service")
            // 关键:设置连接池上限与超时,避免下游慢时把连接耗尽
            .clientConnector(new ReactorClientHttpConnector(
                HttpClient.create()
                    .responseTimeout(Duration.ofSeconds(3))
                    .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 1000)
                    .doOnConnected(conn -> conn
                        .addHandlerLast(new ReadTimeoutHandler(3))
                        .addHandlerLast(new WriteTimeoutHandler(3)))))
            .build();
    }
}

@Service
public class UserReactiveService {

    private final WebClient userWebClient;
    public UserReactiveService(WebClient userWebClient) { this.userWebClient = userWebClient; }

    // 返回 Mono,不阻塞。调用方继续往下编排,等结果回来再处理
    public Mono<UserDTO> getUser(Long id) {
        return userWebClient.get()
            .uri("/user/{id}", id)
            .retrieve()
            // 4xx/5xx 时转成异常,否则 retrieve() 会把错误当正常响应返回
            .onStatus(status -> status.isError(),
                      resp -> resp.bodyToMono(String.class)
                                  .flatMap(body -> Mono.error(
                                      new RuntimeException("下游返回 " + resp.statusCode() + ": " + body))))
            .bodyToMono(UserDTO.class)
            .timeout(Duration.ofSeconds(3))            // 兜底超时,防止无限等待
            .onErrorResume(e -> Mono.just(UserDTO.guest()));  // 降级
    }

    // 并发调用两个下游,再合并结果:总耗时 = max(两个耗时),而不是相加
    public Mono<OrderDetailVO> getOrderDetail(Long userId) {
        Mono<UserDTO>  user  = getUser(userId);
        Mono<List<OrderDTO>> orders = getOrders(userId);
        return Mono.zip(user, orders)
                   .map(t -> new OrderDetailVO(t.getT1(), t.getT2()));
    }

    private Mono<List<OrderDTO>> getOrders(Long userId) {
        return userWebClient.get().uri("/order/{id}", userId)
            .retrieve().bodyToFlux(OrderDTO.class).collectList();
    }
}

函数式端点(RouterFunction)与注解式端点的差别

// ================= 方式一:注解式(和 Spring MVC 写法几乎一样)=================
@RestController
@RequestMapping("/api/user")
public class UserController {

    private final UserReactiveService service;
    public UserController(UserReactiveService service) { this.service = service; }

    // 注意区别:返回类型是 Mono/Flux,而不是普通对象
    // 千万不要在里面 .block()!那等于把响应式又改回了阻塞
    @GetMapping("/{id}")
    public Mono<UserDTO> getUser(@PathVariable Long id) {
        return service.getUser(id);        // 直接返回,框架负责订阅并写出响应
    }

    // Flux 用于流式返回:客户端会分批收到数据(如 SSE)
    @GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public Flux<String> stream() {
        return Flux.interval(Duration.ofSeconds(1)).map(i -> "tick " + i);
    }
}

// ================= 方式二:函数式端点(WebFlux 特有)=================
// 路由和处理器分开写:路由负责"路径 → 处理器",处理器负责业务逻辑
@Configuration
public class UserRouter {

    @Bean
    public RouterFunction<ServerResponse> userRoutes(UserHandler handler) {
        return RouterFunctions.route()
            .GET("/fn/user/{id}", accept(MediaType.APPLICATION_JSON), handler::getUser)
            .POST("/fn/user", handler::createUser)
            .GET("/fn/users", handler::listUsers)
            .build();     // 链式声明,路由表一眼可见
    }
}

@Component
public class UserHandler {

    private final UserReactiveService service;
    public UserHandler(UserReactiveService service) { this.service = service; }

    public Mono<ServerResponse> getUser(ServerRequest request) {
        Long id = Long.valueOf(request.pathVariable("id"));
        return service.getUser(id)
            .flatMap(user -> ServerResponse.ok().bodyValue(user))
            .switchIfEmpty(ServerResponse.notFound().build());
    }

    public Mono<ServerResponse> createUser(ServerRequest request) {
        return request.bodyToMono(UserDTO.class)
            .flatMap(service::save)
            .flatMap(saved -> ServerResponse.created(
                    URI.create("/fn/user/" + saved.getId())).bodyValue(saved));
    }
}

// 对比结论:注解式上手快、和 MVC 一致;函数式路由集中、易做动态组合与测试
// 团队熟悉度优先时选注解式;需要"路由即数据"(比如按配置动态拼装)时选函数式

R2DBC 与响应式数据访问

到这里会遇到一个很尴尬的现实:业务逻辑可以写得很响应式,但 JDBC 是阻塞的。JDBC 规范里没有异步接口,一旦调用 Statement.executeQuery(),线程就必须等。在 WebFlux 里用 JDBC,等于在高速公路上开拖拉机。

方案做法代价 / 适用
JDBC + boundedElastic用 Mono.fromCallable(() -> jdbcCall()) 把阻塞调用丢到弹性线程池改动小、能跑通。但线程池仍会被占满,只是把问题从事件循环挪到了另一个池子。适合过渡期,不适合长期高并发。
R2DBC用响应式的数据库驱动(r2dbc-postgresql、r2dbc-mysql),返回 Mono/Flux真正的非阻塞。但生态较新:部分 ORM 特性缺失(复杂关联、懒加载、二级缓存都没有),动态 SQL 要自己拼或配合 Spring Data R2DBC。
JPA / MyBatis(阻塞)保持原来的数据访问层不能在事件循环里直接用,必须隔离到 boundedElastic。如果数据访问是主要瓶颈,说明这个项目不该用 WebFlux。
混合架构WebFlux 做网关/聚合层(IO 密集),内部服务仍是 MVC + JPA最务实的方案。只在真正需要抗高并发 IO 的那一层用响应式,其余保持简单。

Spring Data R2DBC:接口写法与 JPA 类似,但返回类型是响应式的

// ---------- 1) 实体:不用 JPA 注解,用 Spring Data 的 @Table / @Id ----------
@Table("users")
public class UserEntity {
    @Id
    private Long id;
    private String name;
    private Integer level;
    // getter / setter 省略
}

// ---------- 2) Repository:继承 ReactiveCrudRepository ----------
public interface UserRepository extends ReactiveCrudRepository<UserEntity, Long> {

    // 方法名派生查询,和 JPA 一样;但返回类型是 Mono / Flux
    Flux<UserEntity> findByLevel(Integer level);

    // 自定义 SQL 用 @Query,注意语法是数据库方言
    @Query("SELECT * FROM users WHERE name LIKE :prefix || '%'")
    Flux<UserEntity> searchByNamePrefix(String prefix);
}

// ---------- 3) 配置:连接工厂 + 事务管理 ----------
@Configuration
public class R2dbcConfig {

    @Bean
    public ConnectionFactory connectionFactory() {
        return ConnectionFactories.get(
            "r2dbc:pool:postgresql://localhost:5432/mydb?user=app&password=secret");
        // r2dbc:pool: 前缀会启用连接池,生产必须用池,别直连
    }

    @Bean
    public ReactiveTransactionManager transactionManager(ConnectionFactory cf) {
        return new R2dbcTransactionManager(cf);
    }
}

// ---------- 4) 事务:注意是 @Transactional 的响应式版本 ----------
@Service
public class UserService {

    @Transactional     // 在 WebFlux 下,这个注解会走 ReactiveTransactionManager
    public Mono<UserEntity> createUser(UserEntity user) {
        // 关键:整条链必须在同一个"响应式事务上下文"里
        // 所以用 flatMap 串联,不要在里面 .subscribe(),也不要返回普通对象
        return userRepository.save(user)
            .flatMap(saved -> recordAuditLog(saved))   // 同一事务内,失败一起回滚
            .thenReturn(user);
    }

    private Mono<Void> recordAuditLog(UserEntity user) {
        return auditRepository.save(new AuditLog(user.getId(), "CREATE")).then();
    }
}
响应式事务的三个隐藏规则

规则一:不能在里面 .subscribe()。一旦你手动订阅,事务上下文就断了,那段逻辑跑在事务之外,回滚不了,也容易漏掉异常。整条链必须是一个被框架订阅的 Publisher。

规则二:不能在里面 .block()。在 WebFlux 的请求线程里 block 会直接抛异常(Reactor 会检测),或者更糟——把事件循环卡住,整个应用停止响应。

规则三:不能混用阻塞的数据源。一个 @Transactional 方法里既用 R2DBC 又用 JDBC,事务管理器不一致,回滚行为完全不可预期。要么全响应式,要么全阻塞,并且把阻塞部分明确隔离出去。

什么场景用 WebFlux,什么场景别碰

这一节可能是整章最有价值的部分。WebFlux 不是"更先进的 MVC",它是为特定问题设计的工具。选错了,你会用十倍的复杂度换来更差的性能。

适合 WebFlux别用 WebFlux(继续用 MVC)
IO 密集型、高并发:网关、BFF 聚合层、大量并发调下游 API复杂事务:一个请求里要写多张表、需要 JPA 的高级特性、需要细粒度锁
流式/推送场景:SSE、WebSocket、实时日志、大文件流式下载团队不熟悉响应式:写出来的代码比能解决的问题还多,排查成本高
链路全程非阻塞:下游都是 HTTP + R2DBC + Redis 响应式客户端数据访问以 JDBC/JPA 为主:一进来就是阻塞,绕不过去
需要精细控制背压:比如从上游消费速度远快于下游处理CPU 密集型业务:纯计算任务用线程池更好,响应式只增加开销
Spring Cloud Gateway:它本身就是基于 WebFlux 的中小项目、QPS 不高:MVC 完全够用,响应式属于自找复杂度

一个务实的技术选型判断流程

第一步:这个服务的瓶颈是 IO 还是 CPU?
    如果是 CPU 密集(加解密、图像处理、复杂计算)→ 用 MVC + 线程池,别用 WebFlux
    如果是 IO 密集 → 继续第二步

第二步:峰值并发大概多少?单机需要多少线程才能扛住?
    QPS × 平均响应时间 = 需要的并发线程数
    例:2000 QPS × 300ms = 600 个线程 → MVC 默认 200 线程池已经不够了
    如果算下来需要几百上千线程 → WebFlux 有明确收益
    如果只需要几十个线程 → MVC 完全够,没必要

第三步:整条链路能做到全非阻塞吗?
    数据访问:能不能上 R2DBC?不能的话,阻塞部分占多大比例?
    下游调用:能不能全换成 WebClient?
    如果有超过 30% 的调用是阻塞的 → WebFlux 收益被吃掉大半,慎重

第四步:团队能承担这个学习成本吗?
    响应式的调试、排查、新人上手都是实打实的成本
    答案是否 → 老实选 MVC,把线程池调大 + 加超时熔断,性价比更高

最终建议:微服务架构里最常见的落地方式是
    「网关层 + BFF 聚合层用 WebFlux,业务服务保持 MVC」
    这样收益最大、风险最小

论为什么"响应式"这个词容易误导人

① 它不是说"响应更快"。对于单个请求,WebFlux 的延迟可能比 MVC 还略高一点(多了编排开销)。它的价值在高并发下的资源效率——用更少的线程撑住更多的并发。

② 它也不是"异步编程"。异步只是手段,核心是非阻塞 + 背压:让线程在等待 IO 时不空转,让消费者不会被生产者淹没。理解了这两点,你就能判断一个场景要不要用它。

③ 最实际的一条建议:如果你的系统现在跑得好好的、没有线程池耗尽的告警,那就别动它。引入响应式的正确理由是"我算过账,线程模型撑不住了",而不是"这个技术听起来先进"。架构升级的收益必须能量化,否则就是给团队增加负担。

三个最常见的致命坑

坑一:在响应式链里调用阻塞 API,卡死事件循环。这是 WebFlux 的头号杀手。比如在 map() 里调用 JDBC、Thread.sleep、或者用 RestTemplate 调下游。Netty 的事件循环线程只有 CPU 核数那么多(常见 8 个),你同步阻塞 200ms,8 个线程瞬间全被占住,应用完全停止响应——比 MVC 挂得还彻底。而且它不会抛错,只是"越来越慢直到不动",排查起来很难。解决办法:凡是阻塞调用,一律用 Mono.fromCallable(...).subscribeOn(Schedulers.boundedElastic()) 隔离出去。

坑二:忘记 subscribe,什么都不发生。响应式是"冷"的——没有订阅就没有执行。Mono 构建完了但不订阅,那段代码永远不跑,不报错、不打日志,你就干等。正确做法:① 在 Controller 里直接返回 Mono/Flux,让框架订阅,这是首选;② 如果确实需要手动触发,用 subscribe() 并务必传入错误处理(subscribe(v -> {}, e -> log.error("...", e))),否则异常会被静默吞掉。

坑三:忘了处理错误,异常被吞。找不到 Mono 的错误处理时,Reactor 默认会打印一段"Operator called default onErrorDropped"的日志。很多人不看日志,就以为是"接口没返回但也不报错"。每个异步链的末端都要有兜底:onErrorResume 降级、onErrorMap 转换异常、或者至少 doOnError 打日志。

记
本章小结

① 核心差别:MVC 是一个请求一个线程、阻塞等 IO;WebFlux 是少量事件循环线程、非阻塞调度。

② Mono 是 0 或 1 个元素,Flux 是 0 到 N 个;不订阅就不执行。

③ map 是同步变换,flatMap 是异步并发(乱序),concatMap 是异步串行(保序);zip 配对、merge 交织。

④ subscribeOn 决定源头在哪执行(只有第一个生效),publishOn 决定下游在哪执行(可多次切换)。

⑤ 阻塞调用必须隔离到 boundedElastic();CPU 密集用 parallel()。

⑥ 用 WebClient 代替 RestTemplate;数据访问要么上 R2DBC,要么把阻塞部分隔离出去。

⑦ 选型看三件事:是不是 IO 密集、需不需要几百个并发线程、链路能不能全非阻塞。都不满足就继续用 MVC。

⑧ 最大的坑是在事件循环里阻塞;最常见的坑是忘记 subscribe 或吞掉异常。

章末面试题

面试快答 · WebFlux 与响应式(5 题)

1.(概念题)WebFlux 一定比 Spring MVC 快吗?为什么?

查看答案

答案:不一定,要看场景。低并发下 WebFlux 不但不快,单请求延迟可能还略高(多了编排开销)。高并发 IO 密集型场景下它明显更优:MVC 的线程池被慢下游占满后会排队甚至拒绝请求,而 WebFlux 用少量事件循环线程就能扛住同样的并发,内存占用也低得多。反过来说,如果链路里有阻塞操作(JDBC、RestTemplate),WebFlux 会因为事件循环被卡住而比 MVC 更差。结论:它不是"更快的 MVC",是"针对高并发 IO 场景的另一套模型"。

2.(操作符题)flatMap 和 concatMap 有什么区别?什么时候必须用后者?

查看答案

答案:flatMap 会并发订阅内部 Publisher,结果顺序不保证,吞吐高;concatMap 会串行订阅,前一个完成后才处理下一个,严格保序,吞吐较低。必须用 concatMap 的场景:结果的顺序有业务含义。比如按用户提交顺序写入审计日志、按批次顺序更新同一个余额(并发会导致更新覆盖)、需要保证"先创建再修改"的执行顺序。判断口诀:只要顺序重要,就用 concatMap;顺序不重要且想更快,用 flatMap。

3.(调度题)subscribeOn 和 publishOn 分别影响什么?为什么写了两个 subscribeOn 只有一个生效?

查看答案

答案:subscribeOn 影响订阅信号向上传播时的执行线程,也就是"整条链从哪个线程开始跑"。因为订阅信号只从最下游往上游走一次,所以离源头最近(链上第一个)的 subscribeOn 生效,后面的被覆盖。publishOn 影响数据信号向下游传播时的线程,它只改变自己之后的操作符的执行线程,可以出现多次、分段切换。记忆方法:subscribeOn 管"起点",publishOn 管"中途"。要处理上游的阻塞调用,只能用 subscribeOn。

4.(工程题)一个 WebFlux 服务上线后,压测 QPS 上不去,CPU 不高但响应越来越慢直到超时。最可能的原因?

查看答案

答案:典型的事件循环被阻塞。表现特征就是"CPU 不高(因为线程在等待而不是在计算)但吞吐下降、延迟上涨"。常见来源:在链里调用了 JDBC/MyBatis、用了 RestTemplate 或同步的 HTTP 客户端、Thread.sleep、大文件同步读写、加了解析大量数据的同步 map。排查:打线程 dump 看 reactor-http-nio-* 线程的栈顶是不是卡在某个 JDBC/Socket read 上。修复:把这些阻塞调用用 Mono.fromCallable(...).subscribeOn(Schedulers.boundedElastic()) 隔离,或者换成响应式的客户端(WebClient / R2DBC)。

5.(选型题)你要新起一个后台管理系统,团队 5 人,都没有响应式经验,QPS 峰值 100。选 MVC 还是 WebFlux?

查看答案

答案:选 MVC,毫不犹豫。理由:① QPS 100 用默认 200 线程的 Tomcat 线程池绰绰有余,响应式带来的并发收益为零;② 团队没有响应式经验,写出来的代码会有更多隐藏 bug(忘记 subscribe、吞异常、在链里阻塞),排查成本远高于收益;③ 后台管理系统通常是"复杂事务 + 复杂查询"为主,这正是 WebFlux 最不擅长的场景(JPA 是阻塞的,事务链更容易踩坑);④ 新人上手成本低,招聘也容易。一句话:技术选型要匹配问题规模,不是匹配技术趋势。