WebFlux 与响应式编程:高并发下的另一条路
前面十章讲的都是 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 能多次切换线程。不能。整条链上只有第一个 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 或吞掉异常。
章末面试题
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 是阻塞的,事务链更容易踩坑);④ 新人上手成本低,招聘也容易。一句话:技术选型要匹配问题规模,不是匹配技术趋势。