Reactor 调度器与实战模式
2026/7/16大约 3 分钟
Reactor 调度器与实战模式
一句话概括
调度器控制响应式操作在哪个线程池上执行——选对调度器,是高性能响应式应用的关键。正确拆分 CPU 密集型与 I/O 密集型任务,才能充分发挥非阻塞的优势。
1. 调度器概览
| 调度器 | 线程池 | 用途 |
|---|---|---|
Schedulers.immediate() | 当前线程 | 测试、简单场景 |
Schedulers.single() | 单线程,可复用 | 共享的单线程池 |
Schedulers.parallel() | 固定线程数(= CPU 核数) | CPU 密集型计算 |
Schedulers.boundedElastic() | 有上限的弹性线程池 | 阻塞 I/O 操作 |
选型原则
任务类型是什么?
├─ CPU 密集型(计算、加密、压缩) → Schedulers.parallel()
├─ 阻塞 I/O(数据库、文件、HTTP 调用) → Schedulers.boundedElastic()
├─ 无阻塞操作 → 不需要指定调度器(Netty 的事件循环已处理)
└─ 测试 → Schedulers.immediate()重要:不要在 Reactor 的 Netty 事件循环中执行阻塞操作!这会导致整个响应式应用卡死。一定要用
boundedElastic()隔离阻塞代码。
2. subscribeOn vs publishOn
| 方法 | 作用范围 | 影响位置 |
|---|---|---|
subscribeOn | 整个链的订阅阶段 | 影响上游(从数据源开始的整个链) |
publishOn | 下游操作符 | 只影响之后的操作符 |
Flux.just("a", "b") // main 线程
.map(s -> s.toUpperCase()) // main 线程
.publishOn(Schedulers.parallel()) // 切换线程
.map(s -> "processed: " + s) // parallel 线程池
.subscribeOn(Schedulers.boundedElastic()) // 订阅在 boundedElastic
.subscribe();线程模型可视化
subscribeOn(boundedElastic)
│
▼
Flux.just("a", "b") ← 在 boundedElastic 上订阅
│
map(toUpperCase) ← 在 boundedElastic 上执行
│
publishOn(parallel) ← 切换线程
│
map("processed: …") ← 在 parallel 上执行经验法则:
subscribeOn在链的前端设置一次即可,除非需要改变数据源线程publishOn可以在链中插入多次,每次切换后续操作符的执行线程- 同时使用时,
subscribeOn影响上游,publishOn影响下游
3. 并行处理
将数据分片到多个 CPU 核心并行计算:
Flux.range(1, 10)
.parallel(4) // 拆成 4 条"轨道"
.runOn(Schedulers.parallel()) // 每条轨道分配一个线程
.map(i -> heavyComputation(i)) // 并行执行
.sequential() // 合并回单流
.subscribe(System.out::println);注意: parallel() 只对 CPU 密集型计算有效。I/O 密集型的并行应该用 flatMap + boundedElastic()。
4. 常见陷阱
❌ 在事件循环中执行阻塞操作
// ❌ 错误
Flux.just("/etc/hosts")
.map(path -> Files.readAllBytes(Paths.get(path))) // 阻塞 I/O!卡 Netty 线程
.subscribe();
// ✅ 正确
Flux.just("/etc/hosts")
.subscribeOn(Schedulers.boundedElastic())
.map(path -> Files.readAllBytes(Paths.get(path))) // 在 elastic 线程池执行
.subscribe();❌ parallel 用错地方
// ❌ 错误:parallel 用于 I/O 没有性能提升
flux.parallel(10).runOn(Schedulers.parallel())
.flatMap(url -> httpClient.get(url)) // I/O 操作
// ✅ 正确:I/O 并行用 flatMap + boundedElastic
flux.flatMap(url -> Mono.fromCallable(() -> httpClient.get(url))
.subscribeOn(Schedulers.boundedElastic()))5. 实战模式
模式一:SSE 事件流推送
模拟 Server-Sent Events,每 100ms 推送一个事件:
public Flux<String> sseStream() {
return Flux.interval(Duration.ofMillis(100))
.take(5)
.map(i -> "data: event-" + i + "\n\n");
}输出:
data: event-0
data: event-1
data: event-2
data: event-3
data: event-4模式二:窗口批处理
每 500ms 收集一批数据进行批量处理:
Flux.interval(Duration.ofMillis(100))
.take(20)
.window(Duration.ofMillis(500)) // 500ms 开一个窗口
.flatMap(Flux::collectList); // 窗口内的数据收集为 List模式三:异步日志缓冲区
// 将日志按 100 条一批批量写入
fluxLogs
.buffer(100)
.subscribe(batch -> logWriter.writeBatch(batch));模式四:熔断 + 降级
public Mono<Order> getOrder(String id) {
return orderService.findById(id)
.timeout(Duration.ofMillis(300))
.retry(2)
.onErrorResume(e -> Mono.just(Order.empty(id)));
}关键总结
| 调度器 | 选用时机 |
|---|---|
immediate() | 测试、简单同步场景 |
single() | 需要单线程顺序执行 |
parallel() | CPU 密集型计算 |
boundedElastic() | 阻塞 I/O 操作 |
| 模式 | 核心操作符 |
|---|---|
| SSE 推送 | Flux.interval + map |
| 窗口批处理 | window + flatMap + collectList |
| 异步缓冲 | buffer |
| 熔断降级 | timeout + retry + onErrorResume |
上一篇:Reactor 错误处理与背压
以上是阶段一 Reactor Core 基础的全部内容。接下来进入阶段二:Spring WebFlux MVC 到 Reactive 的迁移。