从同步阻塞到异步非阻塞,用一套代码跑通全链路响应式
去年我们团队接了一个 IoT 数据平台的项目。设备数量从几千台暴涨到几十万台,每台设备每隔几秒就上报一次数据。原来用 Spring MVC 写的接口,扛到 5 万并发就开始频繁超时,线程池打满,GC 频繁,服务器 CPU 倒是没跑满,但大部分线程都在等数据库返回。
后来我们把核心链路改成了 WebFlux + R2DBC,同样的硬件配置,并发吞吐量翻了将近三倍。这不是 WebFlux 比 Spring MVC 快,而是它用更少的资源做了更多的事。
这篇文章不讲太多理论,直接从代码出发。从 Reactor 的 Mono 和 Flux 开始,到 WebFlux 的 Controller 怎么写,再到 R2DBC 怎么连接数据库——全链路走一遍响应式。
一、响应式到底解决了什么问题?
在聊 WebFlux 之前,得先搞清楚一个最基本的问题:我们为什么需要响应式?
传统的 Spring MVC 基于 Servlet 容器(Tomcat、Jetty),每个请求分配一个线程。这个线程从请求进来一直持有到响应返回,中间不管是查数据库、调外部 API,还是读写文件,线程都在那里等着——阻塞。
假设一个请求处理需要 100ms,其中 80ms 花在等数据库返回。在高并发场景下,Tomcat 默认 200 个线程,每秒最多能处理 2000 个请求(200 × 1000ms / 100ms)。线程池满了之后,新的请求只能排队等待。
问题在于:80ms 的等待时间里,线程什么都没干,但资源被占着。CPU 并没有满,但系统已经扛不住了——这就是典型的“线程耗尽”问题。
WebFlux 换了一种思路。它基于 Netty 的事件循环模型,用少量线程处理大量并发请求。请求进来时不分配专属线程,而是注册一个回调,等数据准备好了再通知。线程不阻塞,一直在干活。

这个区别在 IO 密集型场景下尤为明显。WebFlux 官网的数据是:同样 4 核 8GB 的机器,WebFlux 能支撑的并发连接数大约是 Spring MVC 的 2-3 倍。
但要注意:WebFlux 不是银弹。如果你的业务逻辑是 CPU 密集型的(大量计算、加密解密),WebFlux 的优势不明显,甚至因为响应式编程的额外开销反而更慢。它适合的是 IO 密集型场景——大量请求,每个请求都在等外部资源。
二、Reactor:响应式编程的基石
WebFlux 底层依赖 Project Reactor,这是 Java 响应式编程的事实标准,实现了 Reactive Streams 规范。
2.1 Mono 和 Flux
Reactor 提供两个核心类型:Mono 和 Flux。
- Mono:表示 0 或 1 个元素的异步序列。适合单条数据——查一个用户、保存一条记录。
- Flux:表示 0 到 N 个元素的异步序列。适合多条数据——查用户列表、订阅消息流。
// Mono 示例:返回单个用户
Mono<User> findUser(Long id) {
return userRepository.findById(id);
}
// Flux 示例:返回用户列表
Flux<User> listUsers() {
return userRepository.findAll();
}
Mono 和 Flux 是惰性的——你定义了一个 Mono,它不会立即执行,只有当你订阅(subscribe)它的时候,数据流才会真正开始流动。这和 Java 8 的 Stream 有点类似,但 Stream 是同步的,Mono/Flux 是异步的。
2.2 常用的操作符
Reactor 提供了一百多个操作符,最常用的几个:
map:同步转换每个元素
Flux.just(1, 2, 3)
.map(i -> i * 2) // 2, 4, 6
.subscribe(System.out::println);
flatMap:异步转换,每个元素可以返回一个新的 Mono/Flux
Flux.just("user1", "user2")
.flatMap(name -> userService.findByName(name)) // 每个 name 异步查库
.subscribe(user -> System.out.println(user));
flatMap 和 map 的区别是:map 是同步的,flatMap 是异步的。flatMap 里可以执行数据库查询、HTTP 调用等耗时操作,而 map 只适合做简单的数据转换。
filter:过滤符合条件的元素
Flux.range(1, 10)
.filter(i -> i % 2 == 0) // 2, 4, 6, 8, 10
.subscribe(System.out::println);
zip:合并多个 Mono/Flux
Mono<User> userMono = userService.findById(1L);
Mono<Order> orderMono = orderService.findByUserId(1L);
Mono.zip(userMono, orderMono)
.map(tuple -> {
User user = tuple.getT1();
Order order = tuple.getT2();
return new UserOrderResponse(user, order);
});
2.3 背压:生产者和消费者的平衡
背压(Backpressure)是响应式编程的核心机制。简单说就是:消费者告诉生产者“我能吃多少,你慢点生产”。
传统的消息队列里,生产者往队列里写数据,消费者从队列里取。如果生产者太快、消费者太慢,队列就会积压,最终内存爆掉。
Reactor 的解决方式是:消费者订阅时告诉生产者自己的消费能力(request(n)),生产者按照这个速度推送数据。如果消费者处理不过来,可以通过 onBackpressureBuffer()、onBackpressureDrop()、onBackpressureError() 等策略来控制。
// 缓冲策略:积压的数据先存起来,但有限制
Flux.range(1, 1000000)
.onBackpressureBuffer(1000) // 最多缓冲 1000 个
.subscribe();
// 丢弃策略:超出能力的数据直接丢弃
Flux.range(1, 1000000)
.onBackpressureDrop()
.subscribe();
// 错误策略:超出能力直接报错
Flux.range(1, 1000000)
.onBackpressureError()
.subscribe();
实际项目中,背压通常由框架底层处理,开发者不需要显式调用 request(n)。但理解这个概念有助于排查生产环境的问题——如果看到 PendingAcquire 相关的警告,多半是背压没配置好。
三、WebFlux Controller:和 Spring MVC 长得像,但本质不同
WebFlux 的 Controller 写法跟 Spring MVC 几乎一样——@RestController、@GetMapping、@PostMapping 这些注解都能用。但返回类型不一样。
3.1 基础 CRUD
@RestController
@RequestMapping("/api/users")
public class UserController {
private final UserService userService;
public UserController(UserService userService) {
this.userService = userService;
}
@GetMapping("/{id}")
public Mono<User> getUser(@PathVariable Long id) {
return userService.findById(id);
}
@GetMapping
public Flux<User> listUsers() {
return userService.findAll();
}
@PostMapping
public Mono<User> createUser(@RequestBody User user) {
return userService.save(user);
}
@PutMapping("/{id}")
public Mono<User> updateUser(@PathVariable Long id, @RequestBody User user) {
user.setId(id);
return userService.update(user);
}
@DeleteMapping("/{id}")
public Mono<Void> deleteUser(@PathVariable Long id) {
return userService.deleteById(id);
}
}
看到没?除了返回类型从 User 变成了 Mono<User>、从 List<User> 变成了 Flux<User>,其他跟 Spring MVC 一模一样。
但底层的执行机制完全不同。Spring MVC 里,Controller 方法返回时数据已经准备好了;WebFlux 里,返回的是一个尚未执行的数据流,框架会在合适的时机订阅它。
3.2 流式响应:Flux 的杀手锏
WebFlux 真正的优势在于流式响应——数据一边产生一边往外吐,不用等所有数据都准备好了再返回。
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> streamEvents() {
return Flux.interval(Duration.ofSeconds(1))
.map(i -> "Event #" + i)
.take(10); // 只发 10 条
}
前端用 EventSource 或 fetch API 接收,每秒收到一条数据。这种模式在实时通知、日志流、AI 流式对话等场景非常实用。
3.3 错误处理
响应式编程的错误处理跟传统方式不太一样。你不能用 try-catch 包住整个方法——因为 Mono/Flux 是惰性的,异常可能发生在订阅之后。
正确的做法是用操作符处理错误:
@GetMapping("/{id}")
public Mono<User> getUser(@PathVariable Long id) {
return userService.findById(id)
.switchIfEmpty(Mono.error(new UserNotFoundException(id)))
.onErrorResume(UserNotFoundException.class, e ->
Mono.just(new User()) // 返回默认值
)
.doOnError(e -> log.error("查询用户失败: {}", e.getMessage()));
}
常用错误处理操作符:
| 操作符 | 作用 |
|---|---|
onErrorReturn | 发生错误时返回一个默认值 |
onErrorResume | 发生错误时切换到另一个 Mono/Flux |
onErrorMap | 将异常转换成另一种异常 |
doOnError | 发生错误时执行副作用(如日志) |
retry | 发生错误时重试 |
四、R2DBC:响应式的关系数据库连接
Controller 用 WebFlux 改成了异步非阻塞,但如果底层数据库访问还是阻塞的,那整个链路就不算真正的响应式。
R2DBC(Reactive Relational Database Connectivity)就是来解决这个问题的。它提供了非阻塞的数据库驱动,让数据库操作也能融入响应式流。
4.1 引入依赖
<dependencies>
<!-- Spring Boot Starter WebFlux -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
</dependency>
<!-- Spring Data R2DBC -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-r2dbc</artifactId>
</dependency>
<!-- R2DBC MySQL 驱动 -->
<dependency>
<groupId>io.asyncer</groupId>
<artifactId>r2dbc-mysql</artifactId>
<version>1.3.0</version>
</dependency>
</dependencies>
注意:R2DBC 的 MySQL 驱动推荐使用
io.asyncer:r2dbc-mysql,这是原dev.miku:r2dbc-mysql的官方继任者。
4.2 配置文件
spring:
r2dbc:
url: r2dbc:mysql://localhost:3306/testdb
username: root
password: 123456
pool:
initial-size: 10
max-size: 50
max-idle-time: 30m
sql:
init:
mode: always
schema-locations: classpath:schema.sql
R2DBC 的连接池是异步的,和 HikariCP 那种同步连接池完全不同。配置里的 initial-size 和 max-size 控制的是连接池里持有的连接数。
4.3 Repository 层
Spring Data R2DBC 提供了 ReactiveCrudRepository 接口,用法跟 Spring Data JPA 的 JpaRepository 很像,但返回的是 Mono/Flux:
@Repository
public interface UserRepository extends ReactiveCrudRepository<User, Long> {
// 方法名推导:根据邮箱查询
Mono<User> findByEmail(String email);
// 根据状态查询用户列表
Flux<User> findByStatus(Integer status);
// 自定义查询
@Query("SELECT * FROM users WHERE name LIKE CONCAT('%', :name, '%')")
Flux<User> findByNameContaining(@Param("name") String name);
}
实体类的定义和 JPA 类似,但用的是 Spring Data 的 @Table 和 @Id:
@Data
@Table("users")
public class User {
@Id
private Long id;
private String name;
private String email;
private Integer status;
private LocalDateTime createdAt;
}
4.4 Service 层
Service 层把 Repository 的 Mono/Flux 串联起来,形成完整的业务流:
@Service
@Slf4j
public class UserService {
private final UserRepository userRepository;
public UserService(UserRepository userRepository) {
this.userRepository = userRepository;
}
public Mono<User> findById(Long id) {
return userRepository.findById(id)
.switchIfEmpty(Mono.error(new UserNotFoundException(id)))
.doOnNext(user -> log.debug("找到用户: {}", user.getName()));
}
public Flux<User> findAll() {
return userRepository.findAll()
.doOnComplete(() -> log.debug("查询所有用户完成"));
}
public Mono<User> save(User user) {
user.setCreatedAt(LocalDateTime.now());
return userRepository.save(user)
.doOnSuccess(saved -> log.info("用户创建成功: id={}", saved.getId()));
}
public Mono<User> update(User user) {
return userRepository.findById(user.getId())
.flatMap(existing -> {
existing.setName(user.getName());
existing.setEmail(user.getEmail());
existing.setStatus(user.getStatus());
return userRepository.save(existing);
})
.switchIfEmpty(Mono.error(new UserNotFoundException(user.getId())));
}
public Mono<Void> deleteById(Long id) {
return userRepository.deleteById(id)
.doOnSuccess(v -> log.info("用户删除成功: id={}", id));
}
}
flatMap 在这里的作用是:先查询用户是否存在,存在则更新,不存在则报错。因为 flatMap 里可以返回一个新的 Mono,适合做这种“先查后改”的操作。
4.5 完整的请求处理流程

从 Controller 到 DB,整条链路都是非阻塞的。没有哪个环节会占着线程等结果。
五、测试:怎么写响应式单元的测试
响应式代码的测试跟传统测试不太一样,不能直接用 assertEquals 去验证 Mono 里的值——因为 Mono 是异步的,测试方法返回时数据可能还没到。
Spring 提供了 StepVerifier 专门用来测试响应式流:
@SpringBootTest
@Slf4j
public class UserServiceTest {
@Autowired
private UserService userService;
@Test
void testFindById() {
// 准备数据
User user = new User();
user.setName("测试用户");
user.setEmail("test@example.com");
Mono<User> savedMono = userService.save(user);
// 用 StepVerifier 验证
StepVerifier.create(savedMono)
.expectNextMatches(saved -> saved.getId() != null)
.verifyComplete();
}
@Test
void testFindByIdNotFound() {
StepVerifier.create(userService.findById(99999L))
.expectError(UserNotFoundException.class)
.verify();
}
@Test
void testFindAll() {
StepVerifier.create(userService.findAll())
.expectNextCount(3) // 期望至少 3 条
.verifyComplete();
}
}
StepVerifier 会订阅 Mono/Flux,等待数据流完成,然后验证每个元素是否符合预期。它支持多种验证方式——expectNext 验证下一个元素、expectNextCount 验证数量、expectError 验证异常。
六、什么时候用 WebFlux?
说了这么多,最后得回到一个现实问题:我的项目该不该上 WebFlux?
适合的场景:
- IO 密集型:大量请求需要查数据库、调外部 API、读写文件。这类场景 WebFlux 的优势最明显。
- 高并发连接:需要维持大量长连接(WebSocket、SSE),WebFlux 的线程模型比 Spring MVC 高效得多。
- 流式数据处理:需要一边产生数据一边返回,比如实时日志、AI 流式对话。
- 网关/代理:Spring Cloud Gateway 本身就是基于 WebFlux 构建的,天然适合做流量入口。
不适合的场景:
- CPU 密集型:大量计算、加密解密、图片处理。这类场景用 WebFlux 反而因为响应式框架的额外开销而更慢。
- 团队不熟悉响应式:响应式编程的学习曲线比想象中陡峭。
flatMap和map的区别、背压的处理、错误传播的方式——这些都需要时间适应。 - 现有代码改造成本高:如果整个项目都是同步代码,强行改成 WebFlux 需要重写 DAO 层、Service 层、Controller 层,改动量非常大。对于这类项目,不如考虑在网关层用 WebFlux,内部服务保持同步。
一个折中的方案:网关层用 Spring Cloud Gateway(基于 WebFlux),处理鉴权、限流、路由;业务服务继续用 Spring MVC。这样既能享受 WebFlux 在高并发入口的优势,又不用大规模改造现有业务代码。
七、总结
WebFlux 不是用来替代 Spring MVC 的,它们是两种不同场景下的选择。
从 Reactor 的 Mono/Flux 到 WebFlux 的异步 Controller,再到 R2DBC 的非阻塞数据库访问——全链路响应式让 Java 后端在面对高并发 IO 场景时有了新的解法。它用更少的线程资源支撑更多的并发连接,在 IoT、实时数据、API 网关等场景下优势明显。
但响应式编程的代价是复杂度。flatMap 嵌套多了可读性会下降,调试时堆栈信息比同步代码难追,团队需要花时间适应新的编程范式。
选不选 WebFlux,取决于你的业务场景和团队状况。如果两者都满足,值得一试。
系列拓展阅读
- 《Spring Cloud Gateway 深度实战:路由、限流、熔断与统一鉴权》 —— WebFlux 在网关层的典型应用
- 《Spring Boot 3.5 AOT 与 GraalVM 原生镜像:启动时间从 3 秒到 0.1 秒》 —— WebFlux 应用编译成原生镜像的注意事项
- 《JUnit 5 + Testcontainers:Java 微服务集成测试最佳实践》 —— 响应式应用的集成测试方案
- 《Spring Data JPA 与 MyBatis 终极对决:Java 数据访问层选型指南》 —— R2DBC 与传统 ORM 的对比
参考文献
- Spring WebFlux Reference Documentation. https://docs.spring.io/spring-framework/reference/web/webflux.html
- Project Reactor Documentation. https://projectreactor.io/docs
- R2DBC Specification. https://r2dbc.io/
- Spring Data R2DBC Reference. https://docs.spring.io/spring-data/r2dbc/reference/
- “万字长文:Spring WebFlux + Project Reactor 配置全解析.” 腾讯云, 2026
- “WebFlux背压策略与Observer模式线程优化实战.” OSC, 2026
- “Java反应式编程ProjectReactor背压问题及生产环境爆炸原因分析.” OSC, 2026
- “SpringBoot3x响应式编程实践指南.” OSC, 2026
- “Reactor与RxJava背压机制实战指南.” OSC, 2026
- “SSE流式响应:从Reactor Flux到生产级AI聊天的工程实践.” CSDN, 2026









