RxJava 4.0 中 cached()、virtual() 与 computation() 三大标准调度器怎么选?
【免费下载链接】RxJavaRxJava – Reactive Extensions for the JVM – a library for composing asynchronous and event-based programs using observable sequences for the Java VM.项目地址: https://gitcode.com/gh_mirrors/rx/RxJava
在 RxJava 4.0 中把阻塞调用移出订阅线程时,必须决定使用哪个标准Scheduler。这个决定绕不开,因为 4.0 中传统的Schedulers.io()已被标记为 API deprecated 并内部委托给Schedulers.cached(),而 快速迁移指南明确要求你在这些io()的旧调用位置逐一决定:是照旧用cached(),还是换用新的virtual()。CPU 密集型任务则属于computation()的范畴。本文基于仓库文档,给出这三个调度器的适用场景、可配置的系统属性,以及跑通和验证每条路径的方法。
准备条件
- Java 版本:RxJava 4.0 是原生 Java 26 实现(见 README 中的版本说明),而
virtual()依赖 Java 21 标准化、4.0 直接集成的虚拟线程基础设施。 - 引入依赖:以 Gradle 为例,把坐标前缀换成 4.x 的
io.reactivex.rxjava4(README 的 Getting started 一节):
implementation "io.reactivex.rxjava4:rxjava:4.x.y"- 包名变化:4.0 的组件位于
io.reactivex.rxjava4,基础类型位于io.reactivex.rxjava4.core。从 3.x 迁移时按 迁移指南把io.reactivex.rxjava3系列导入替换为io.reactivex.rxjava4系列即可。
三个调度器的定位
README 的 Schedulers 一节和 Schedulers 的 Javadoc 对三者的分工描述如下:
| 调度器 | 适用工作 | 底层实现 |
|---|---|---|
Schedulers.computation() | 计算密集型任务、事件循环、处理回调;大多数异步算子默认使用它 | 固定数量(默认等于可用处理器数)的单线程ScheduledExecutorService实例 |
Schedulers.cached() | I/O 类或阻塞操作,即原Schedulers.io()的改名 | 单线程ScheduledExecutorService实例池,worker 线程数量无界 |
Schedulers.virtual() | I/O 类或阻塞的 scatter-gather 操作,以顺序化的方式运行 | 每个 Worker 使用 Java 标准的Executors.newVirtualThreadPerTaskExecutor() |
关键取舍依据来自文档原文:
cached()可能创建无界的 worker 线程,导致系统减速甚至OutOfMemoryError;README 对它直接标注了 "Can exhaust system resources!"。virtual()的价值在于 "Helps with the issues around unboundedness ofSchedulers.cached()",即缓解cached()的无界问题。What's-different-in-4.0.md 的说法是:通过虚拟线程可以用阻塞 API 发起成千上万次 web 或数据库调用而不耗尽系统资源。computation()不建议跑阻塞或 IO 密集型工作;cached()不建议跑计算型工作,两者在 Javadoc 里互相指向对方。
执行路径一:阻塞 IO 改用 cached() 或 virtual()
迁移场景下最常见的任务,是把subscribeOn(Schedulers.io())替换成下面两者之一。两者用法位置相同,只是调度器不同。
用virtual()运行阻塞调用(写法来自 What's-different-in-4.0.md 中Schedulers.virtual()一节,文档示例里的yourBlockingCallHere()需替换为你自己的阻塞调用):
import io.reactivex.rxjava4.core.Flowable; import io.reactivex.rxjava4.schedulers.Schedulers; Flowable.fromCallable(() -> { Thread.sleep(1000); // 占位:替换为真实的阻塞调用(文档示例用 yourBlockingCallHere()) return "Done"; }) .subscribeOn(Schedulers.virtual()) .subscribe(System.out::println);用cached()运行阻塞调用(subscribeOn+ 阻塞 callable 的组合方式来自 README 的示例,该示例同时展示了下面要说明的 daemon 线程注意点):
Flowable.fromCallable(() -> { Thread.sleep(1000); // 模拟耗时的阻塞计算 return "Done"; }) .subscribeOn(Schedulers.cached()) .subscribe(System.out::println, Throwable::printStackTrace); Thread.sleep(2000); // 保持 main 线程存活,见下文两条路径都依赖一个 README 明确指出的细节:RxJava 的标准Scheduler运行在daemon 线程上,main 线程退出后它们全部停止,后台计算可能永远不会执行。因此在控制台演示或短生命周期示例中,需要让 main 线程多活一会儿(如上例的Thread.sleep(2000))。
io()没有删除,只是为了不让存量代码库出现找不到方法的编译错误。它@Deprecated(since = "4.0.0"),Javadoc原话是 "please stop using this",并指向cached()或virtual()。
执行路径二:CPU 密集工作走 computation()
README 给出的并发示例可以直接运行,用于把处理搬到computation():
import io.reactivex.rxjava4.core.Flowable; import io.reactivex.rxjava4.schedulers.Schedulers; Flowable.range(1, 10) .observeOn(Schedulers.computation()) .map(v -> v * v) .blockingSubscribe(System.out::println);注意 README 对这个示例的说明:v -> v * v并不会并行执行,1 到 10 会在同一个 computation 线程上按顺序逐个处理。若需要并行处理,文档给出的是flatMap+ 每路subscribeOn的写法,这里只作了解即可,不属于本文的选型主路径。
可调系统属性
以下属性必须在Schedulers类被引用之前通过System.getProperty方式设置(见 Schedulers Javadoc 与 What's-different-in-4.0.md 的 Cached scheduler 一节;后者的前缀由rx3.io*改名为rxjava4.cached*):
rxjava4.cached-keep-alive-time(long):cached()worker 的 keep-alive 时间,默认值见CachedScheduler.KEEP_ALIVE_TIME_DEFAULT。rxjava4.cached-priority(int):cached()的线程优先级,默认Thread.NORM_PRIORITY。rxjava4.cached-scheduled-release(boolean):设为true时把 worker 释放模式从默认的 eager 改为 scheduled。Javadoc 说明了两模式的取舍——eager(默认)下 worker 立即回池、可更快复用,但如果当前任务不响应中断,复用可能导致延迟或死锁;scheduled 下 worker 等当前任务结束才回池,能减少提前复用,但可能产生过多的底层 worker。rxjava4.computation-threads(int):computation()的线程数,默认为可用 CPU 数。rxjava4.computation-priority(int):computation()的线程优先级,默认Thread.NORM_PRIORITY。
验证与检查
- 看控制台输出:README 示例的成功标志就是能在控制台看到 flow 的输出(配合上文
Thread.sleep保持 main 存活)。文档没有给出固定成功日志,输出内容即你的订阅回调打印的内容。 - 阻塞保护开关:RxJavaPlugins.setFailOnNonBlockingScheduler 设为
true后,在computation()(以及single())这类非阻塞调度器上执行阻塞算子会抛出IllegalStateException。这是一个现成的手段:开启它来捕获被误放到computation()上的阻塞代码。注意该设置在插件被lockdown()后不可再修改。 - 未处理错误的去向:三个调度器上未处理的错误都会交给该调度器线程的
Thread.UncaughtExceptionHandler(SchedulersJavadoc 对computation()、cached()均如此说明),排查挂死的流程时可以从这里入手。
限制与边界
cached()的 worker 线程无界, casual 使用或实现算子时必须通过Worker.dispose()释放,否则可能导致系统减速或OutOfMemoryError(Javadoc 原文)。virtual()调试受限:What's-different-in-4.0.md 指出部分 IDE(如 Eclipse)无法正确对虚拟线程代码单步,断点可能失效或挂起 IDE;文档给出的替代做法是在单元测试或调试时改用传统调度器运行同一段代码。- 标准的
Schedulers.io()已 deprecated,代码里见到它时应视为待处理项,而不是默认选择。 - 应用退出前可调用
Schedulers.shutdown()关闭标准调度器(幂等且线程安全);通过RxJavaPlugins.createXxxScheduler(ThreadFactory)创建的自定义实例则需要手动调用Scheduler.shutdown()才能让 JVM 正常退出。
完整的调度器行为说明见 Scheduler.md 指向的 Scheduler 文档、Schedulers 的 Javadoc,以及 What's-different-in-4.0.md 中 "New schedulers" 与 "Cached scheduler" 两节。
【免费下载链接】RxJavaRxJava – Reactive Extensions for the JVM – a library for composing asynchronous and event-based programs using observable sequences for the Java VM.项目地址: https://gitcode.com/gh_mirrors/rx/RxJava
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考