虚拟线程出现后,「是否还需要响应式编程」成为热门话题。答案是:在高吞吐流式场景下,Reactor 依然不可替代;在普通 CRUD 服务中,虚拟线程是更简单的选择。本文聚焦 Project Reactor 的核心概念和实战模式。
Mono 与 Flux:两个核心类型
Project Reactor 基于 Reactive Streams 规范,提供两个发布者类型:
- Mono<T>:0 或 1 个元素的异步序列(类似 Optional + Future)
- Flux<T>:0 到 N 个元素的异步序列(类似 Stream + Future)
// 查询单个用户
Mono<User> findById(String id) {
return userRepository.findById(id)
.switchIfEmpty(Mono.error(new UserNotFoundException(id)));
}
// 流式返回大量日志
Flux<LogEntry> streamLogs(Instant since) {
return logRepository.findByTimestampAfter(since);
}
操作符精选
Reactor 有 200+ 操作符,日常开发中最高频的是这些:
- map / flatMap:同步/异步转换元素
- filter / take / skip:过滤和截取
- zip / merge:合并多个流
- onErrorResume / onErrorReturn:错误恢复
- timeout / retry:超时与重试
- publishOn / subscribeOn:线程调度
flatMap 的并发度(concurrency)是最容易被忽视的参数——默认值可能导致下游服务被打爆。
Spring WebFlux 集成
WebFlux 让 Controller 直接返回 Mono/Flux,底层使用 Netty 事件循环:
@RestController
@RequestMapping("/api/users")
public class UserController {
@GetMapping("/{id}")
public Mono<UserDto> getUser(@PathVariable String id) {
return userService.findById(id)
.map(UserDto::from)
.timeout(Duration.ofSeconds(3));
}
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<UserDto> streamUsers() {
return userService.findAll()
.map(UserDto::from)
.delayElements(Duration.ofMillis(100));
}
}
何时选 Reactor,何时选虚拟线程
| 场景 | 推荐方案 | 原因 |
|---|---|---|
| 普通 REST CRUD | 虚拟线程 + 阻塞式 | 代码简单,性能够用 |
| SSE / WebSocket 流 | WebFlux + Flux | 原生背压支持 |
| 网关 / BFF 聚合 | WebFlux + flatMap | 并发调用多个下游 |
| 消息流处理 | Reactor Kafka | 端到端响应式管道 |
| CPU 密集计算 | 虚拟线程 + 线程池隔离 | Reactor 事件循环不适合 CPU 密集 |
常见陷阱
- 阻塞调用混入响应式链:用
subscribeOn(Schedulers.boundedElastic())包裹,或改用 R2DBC 替代 JDBC - 忘记 subscribe:Reactor 是惰性的,不 subscribe 什么都不会发生
- 异常被吞:务必在链中添加
doOnError日志,或使用Hooks.onErrorDropped监控 - 测试困难:用
StepVerifier做单元测试,比 block() 更可靠
小结
响应式编程不是银弹,但在流式数据、高并发 IO 聚合、背压控制场景下,Project Reactor 提供了比回调地狱更优雅的抽象。理解 Mono/Flux 语义,选对使用场景,比全面响应式化更重要。