响应式编程带来的最大改变,不在于语法上的 Mono 和 Flux,而在于整个请求处理链路从阻塞变为非阻塞。在 Spring WebFlux 中开发 Web 接口,本质上就是将入站请求、业务处理、数据库访问、外部服务调用全部串联成一个异步非阻塞的流,用声明式的方式描述数据如何流动与变换。
25.3.1 两种编程模型的选择
Spring WebFlux 提供了两种编写接口的方式:
- 注解式控制器(Annotated Controllers):与 Spring MVC 使用相同的
@RestController、@GetMapping等注解,唯一的区别是方法返回值变成了Mono<T>或Flux<T>。这是绝大多数团队的默认选择,迁移成本极低。 - 函数式端点(Functional Endpoints):使用
RouterFunction和HandlerFunction以纯 Lambda 风格定义路由和处理逻辑,无任何注解,适合追求极致轻量与函数式风格的场景。
考虑到实际项目中最常见的仍是注解式控制器,本节将以它为主线展开,同时简要介绍函数式端点以便读者在需要时快速上手。
25.3.2 注解式响应式控制器
一个典型的响应式 REST 接口看起来和传统 MVC 接口几乎一样,只是返回值包装在 Mono 或 Flux 中:
@RestController
@RequestMapping("/api/users")
public class UserController {
private final UserRepository userRepository;
public UserController(UserRepository userRepository) {
this.userRepository = userRepository;
}
@GetMapping("/{id}")
public Mono<User> getUser(@PathVariable String id) {
return this.userRepository.findById(id);
}
@GetMapping
public Flux<User> listUsers(@RequestParam(defaultValue = "0") int page,
@RequestParam(defaultValue = "10") int size) {
return this.userRepository.findAll()
.skip((long) page * size)
.take(size);
}
@PostMapping
@ResponseStatus(HttpStatus.CREATED)
public Mono<User> createUser(@RequestBody Mono<User> userMono) {
return userMono.flatMap(this.userRepository::save);
}
@DeleteMapping("/{id}")
public Mono<Void> deleteUser(@PathVariable String id) {
return this.userRepository.deleteById(id);
}
}
关键点解释:
Mono<User>表示“未来某个时刻会就绪一个 User 或者空”。框架订阅这个Mono,并在值就绪时将其序列化为 JSON 写入 HTTP 响应。Flux<User>表示“一个可能包含 0 到 N 个 User 的流”。框架会逐个将这些 User 序列化,并以分块传输的方式发送给客户端,客户端可以边接收边处理。- 请求体可以直接映射为
Mono<T>或Flux<T>,让入站的反序列化也变成非阻塞过程。比如@RequestBody Mono<User>,Spring 会异步地读取和解析 HTTP 请求体,解析完成后发出一个User对象。 userRepository通常是一个响应式的数据访问接口(如ReactiveCrudRepository),其方法本身就返回Mono或Flux,因此可以在控制器中直接返回,形成一个端到端的非阻塞链路。
不需要手动订阅。 这是很多从传统编程转入响应式时容易忽视的一点:控制器的返回值只是一种描述,WebFlux 框架本身会作为最终的订阅者,负责将数据推送到网络层。开发者的职责只是将上游的响应式类型进行组合、变换,最后扔给框架,让整条链路保持非阻塞。
25.3.3 服务端推送事件(Server-Sent Events)
对于需要持续推送数据的场景(如股票行情、进度条更新),Flux 可以进一步配合 produces 属性输出 SSE 流:
@GetMapping(value = "/{id}/events", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<Event> streamEvents(@PathVariable String id) {
return this.eventService.eventsForUser(id)
.delayElements(Duration.ofSeconds(1)); // 模拟每秒推送一条
}
此时浏览器或其他客户端可以使用 EventSource API 持续接收事件。如果将 MediaType 改为 APPLICATION_NDJSON_VALUE,则变为 NDJSON(换行分隔 JSON)流,同样适合流式数据处理。
25.3.4 统一错误处理
在响应式控制器中,异常会通过 Mono.error() 或 Flux.error() 传播,而不是被传统的 @ExceptionHandler 直接捕获。Spring WebFlux 支持在响应式上下文中使用 @RestControllerAdvice,只需确保方法返回的是 Mono 或 Flux 封装的结果:
@RestControllerAdvice
public class GlobalExceptionHandler {
@ExceptionHandler(UserNotFoundException.class)
@ResponseStatus(HttpStatus.NOT_FOUND)
public Mono<ErrorResponse> handleNotFound(UserNotFoundException ex) {
return Mono.just(new ErrorResponse(ex.getMessage(), 404));
}
}
对于需要在流中处理单条记录异常的场景,可以在 Flux 上使用 onErrorContinue 或 onErrorResume 等操作符,避免整个流因一条数据而终止。
25.3.5 请求体与参数校验
注解式响应式控制器同样支持标准的 Bean Validation。只需在方法参数上加上 @Valid,并将请求体包装为 Mono<T>,框架会自动在反序列化后执行校验:
@PostMapping
public Mono<User> createUser(@Valid @RequestBody Mono<User> userMono) {
return userMono.flatMap(userService::create);
}
如果校验失败,会抛出一个 WebExchangeBindException(对应 MVC 中的 MethodArgumentNotValidException),可以在全局异常处理器中统一格式化返回的错误信息。
25.3.6 函数式端点(Router Functions)
对于希望摆脱注解、完全用代码表述路由的团队,可以定义 RouterFunction 和 HandlerFunction:
@Configuration
public class UserRouter {
@Bean
public RouterFunction<ServerResponse> userRoutes(UserHandler handler) {
return RouterFunctions
.route(GET("/api/users/{id}").and(accept(MediaType.APPLICATION_JSON)), handler::getUser)
.andRoute(GET("/api/users").and(accept(MediaType.APPLICATION_JSON)), handler::listUsers)
.andRoute(POST("/api/users").and(accept(MediaType.APPLICATION_JSON)), handler::createUser);
}
}
处理器则是普通的 Bean:
@Component
public class UserHandler {
private final UserRepository userRepository;
public UserHandler(UserRepository userRepository) {
this.userRepository = userRepository;
}
public Mono<ServerResponse> getUser(ServerRequest request) {
String id = request.pathVariable("id");
return this.userRepository.findById(id)
.flatMap(user -> ServerResponse.ok().bodyValue(user))
.switchIfEmpty(ServerResponse.notFound().build());
}
public Mono<ServerResponse> listUsers(ServerRequest request) {
Flux<User> users = this.userRepository.findAll();
return ServerResponse.ok().body(users, User.class);
}
public Mono<ServerResponse> createUser(ServerRequest request) {
return request.bodyToMono(User.class)
.flatMap(userRepository::save)
.flatMap(saved -> ServerResponse.created(URI.create("/api/users/" + saved.getId()))
.bodyValue(saved));
}
}
函数式端点的优势在于极致的显式化和可组合性。你可以将路由分组、嵌套、中间件过滤,全部通过代码完成,不需要依赖注解扫描。对于某些需要动态生成路由、或对反射机制敏感的项目,这是一种非常干净的选择。
25.3.7 背压与性能考量
响应式 Web 接口并非总是更快,它的真正威力体现在缓慢的下游消费者或高并发 I/O 密集型场景。
假设你的接口从数据库读取 10 万条记录,客户端由于网络或处理能力限制只能每秒消费 100 条。在传统的阻塞式 List<User> 返回方式下,所有数据必须一次性加载到内存并序列化完成,然后才开始传输;而在 Flux<User> 方式下,Spring WebFlux 会感知 TCP 写缓冲区的容量,按需向数据库请求下一批行,这就是背压。内存消耗由连接数和缓冲区大小决定,与数据总量无关。
实际开发中,为了利用这一特性,必须保证整条链路的响应式:
- Web 层:Spring WebFlux + Netty
- 数据访问:Spring Data R2DBC / Reactive MongoDB / Reactive Redis
- HTTP 客户端:
WebClient - 外部服务调用全程返回
Mono/Flux
只要有任何一个环节阻塞(例如调用了传统的 JDBC),整个流的非阻塞特性就会被破坏,最终退化为在单独的线程池中运行,性能收益会大幅折扣。
25.3.8 实用的测试方式
@WebFluxTest 提供了对响应式控制器的切片测试支持:
@WebFluxTest(UserController.class)
class UserControllerTest {
@Autowired
private WebTestClient webTestClient;
@MockBean
private UserRepository userRepository;
@Test
void shouldReturnUserWhenExists() {
User user = new User("1", "张三");
when(userRepository.findById("1")).thenReturn(Mono.just(user));
webTestClient.get().uri("/api/users/1")
.accept(MediaType.APPLICATION_JSON)
.exchange()
.expectStatus().isOk()
.expectBody(User.class)
.isEqualTo(user);
}
@Test
void shouldReturn404WhenNotFound() {
when(userRepository.findById("99")).thenReturn(Mono.empty());
webTestClient.get().uri("/api/users/99")
.exchange()
.expectStatus().isNotFound();
}
}
WebTestClient 是一个反应式的 HTTP 客户端,可以直接针对RouterFunction或运行中的服务器进行测试,无需启动 Servlet 容器。
至此,响应式 Web 接口开发的完整路径已经从原理到实践落地。它的本质是异步流的拼接,注解式方式和函数式方式都只是描述这一流的手段。选择哪一种取决于团队的编码习惯,但无论哪种,保证链路完全非阻塞才是发挥 WebFlux 价值的前提。