本节摘要:
Scheduler不是线程池别名,而是时间语义仲裁器:立即、延迟、周期、重试共享同一抽象。Reactor 提供boundedElastic(阻塞 I/O)、parallel(CPU)、single(严格顺序)与测试用VirtualTimeScheduler。
subscribeOn 与 publishOn风控引擎同时接入交易流水(高吞吐)、画像更新(中频)、黑名单推送(偶发)。若用三线程池 + 锁,隐含假设是资源充裕、失败例外。响应式问:若下游永远慢于上游,系统如何自我校准?答案在 Scheduler + 背压,而非更大线程池。
FRP 将信号建模为 ( s: \mathbb{T} \to \mathbb{V} )——时间域到值域的函数,使 delay、sample、timeout 成为合法组合。
时间的显式化带来一个直接结果:时间操作符成为可组合的一等公民。命令式里实现「500ms 防抖」要在回调里维护定时器状态;响应式里就是 debounce(Duration.ofMillis(500)) 一个算子,它内部管理定时器、取消旧定时器、串行化信号,这些复杂性被封装在算子语义里。delay、sample、timeout、interval、take(duration) 全部基于同一个时间抽象——Scheduler。这就是「时间是一等公民」的工程含义:时间不是散落在代码里的 Thread.sleep,而是可声明、可测试、可组合的流属性。
| Scheduler | 用途 | 注意 |
|---|---|---|
Schedulers.immediate() |
当前线程 | 无切换 |
Schedulers.boundedElastic() |
阻塞 I/O | 有界线程,防爆炸 |
Schedulers.parallel() |
CPU 计算 | 默认 CPU 核数 |
Schedulers.single() |
全局顺序 | 单线程串行 |
Mono.fromCallable(() -> jdbcQuery()) .subscribeOn(Schedulers.boundedElastic()) .flatMap(row -> webClient.post().bodyToMono(Void.class));
subscribeOn:影响源订阅发生的线程publishOn:影响后续算子执行线程嵌套错误会导致背压信号丢失(第 6 章反模式「背压盲区」)。
选型规则可以浓缩成一句口诀:阻塞的活儿给 boundedElastic,纯计算的活儿给 parallel,必须严格串行的给 single,不想切换就用 immediate。 特别要记住 boundedElastic 的名字——它默认有 10 倍 CPU 核数的线程上限,专为「不得不调阻塞代码」设计,防止你把 EventLoop 全部阻塞死。而 parallel 的数量等于 CPU 核数,适合无阻塞的 CPU 密集变换。
// subscribeOn 与 publishOn 的差异演示 Flux.fromIterable(list) // 在调用线程创建 .subscribeOn(Schedulers.boundedElastic()) // 源与上游在弹性线程 .map(this::parse) // 仍在 boundedElastic .publishOn(Schedulers.parallel()) // 从这往下的算子切到 parallel .map(this::score) // 在 parallel 执行 .subscribe(); // 下游消费在 parallel
关键在切换点:subscribeOn 只影响它上游(源侧)的执行线程,publishOn 影响它下游(后续算子)的执行线程。上例中 parse 跑在 boundedElastic(因为还没遇到 publishOn),score 跑在 parallel(publishOn 之后的链段)。理解这个「切刀位置」后,你就能精确控制每段代码跑在哪种线程上——这是响应式性能调优的基本功,也是第六章调试时定位「哪一段在阻塞」的依据。
StepVerifier 配合 VirtualTimeScheduler 可将 Flux.interval(Duration.ofHours(1)) 在测试中瞬间推进,无需真等一小时(SOURCE 第六章测试策略)。
RxJS 对照:asyncScheduler、queueScheduler、animationFrameScheduler 分别服务微任务、队列与浏览器帧;observeOn/subscribeOn 语义与 Reactor 的 publishOn/subscribeOn 类似。
// 虚拟时间:把 1 小时的间隔在毫秒级测试中推进 StepVerifier.withVirtualTime(() -> Flux.interval(Duration.ofHours(1)).take(3)) .expectSubscription() .thenAwait(Duration.ofHours(3)) // 虚拟推进 3 小时 .expectNext(0L, 1L, 2L) .verifyComplete();
虚拟时间的原理是替换 Scheduler:测试框架把全局时间源换成可手动推进的虚拟时钟,所有 delay/interval/timeout 的定时都基于这个虚拟时钟,于是 thenAwait(Duration.ofHours(3)) 瞬间完成。这带来两个决定性优势:测试不再依赖真实等待(秒级甚至零耗时);测试具备确定性——不会因为机器负载而偶发超时(flaky)。代价是你必须保证被测代码的时间操作都经过 Scheduler 抽象,若有人偷偷 Thread.sleep,虚拟时间就失效了——这也是第六章把它列入反模式的原因。

调度器不是孤立的工具,它直接支撑 Manifesto 的 Elstic 与 Responsive:
timeout(Duration) 依赖调度器计量时间,超时降级是「及时响应」的执行保证;boundedElastic 的线程上限就是「可伸缩但有边界」的运行时形态——负载高时线程数顶到上限而不是无限涨,配合背压(2.3)形成可控的队列水位;所以记住:Scheduler 是时间与线程的仲裁层,它决定了一段流「在哪个线程、何时、以什么节奏」执行。它是响应式系统的「时钟 + 执行器」二合一。
时间作为一等公民的另一个体现,是存在一整族时间操作符:delayElements、interval、timeout、sample、take(duration)、debounce、throttleFirst。它们共享同一个时间抽象,因而可以自由组合:
| 操作符 | 语义 | 典型场景 |
|---|---|---|
delayElements(d) |
每个元素延后 d | 平滑推送节奏 |
interval(d) |
每 d 发射递增序号 | 心跳、轮询 |
timeout(d) |
超时未发射则报错 | 响应预算 |
sample(d) |
每 d 取最新一次 | 节流采样 |
take(d) |
取前 d 时长内元素 | 限时窗口 |
debounce(d) |
静默 d 后发射 | 输入防抖 |
// 组合示例:限时窗口 + 防抖 Flux.merge( clickStream.debounce(Duration.ofMillis(300)), // 防抖 timerStream.take(Duration.ofSeconds(30))) // 30 秒后终止 .subscribe(handle);
// RxJS 的时间操作符与 Reactor 对应 import { fromEvent, debounceTime, throttleTime, sampleTime } from 'rxjs'; const clicks = fromEvent(document, 'click'); clicks.pipe(debounceTime(300)).subscribe(...); // 静默 300ms 后 clicks.pipe(throttleTime(1000)).subscribe(...); // 1 秒最多一次 clicks.pipe(sampleTime(2000)).subscribe(...); // 每 2 秒取最新
这些操作符的存在说明一个深刻的变化:在命令式世界里,时间操纵散落在定时器、回调与 sleep 里;在响应式世界里,时间被压缩进操作符的参数,成为可以声明、测试与组合的语言构件。 2.2 的全部分量——Scheduler、虚拟时间、时间操作符——都指向同一个结论:时间是响应式的第一类公民,而调度器是它的运行时代表。
boundedElastic下一节专讲背压:下游如何用
request(n)反向控流。