十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

RxJava 4.0 中 cached()、virtual() 与 computation() 三大标准调度器怎么选?

RxJava 4.0 中 cached()、virtual() 与 computation() 三大标准调度器怎么选? 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.rxjava4README 的 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 线程导致系统减速甚至OutOfMemoryErrorREADME 对它直接标注了 Can exhaust system resources!。virtual()的价值在于 Helps with the issues around unboundedness ofSchedulers.cached()即缓解cached()的无界问题。Whats-different-in-4.0.md 的说法是通过虚拟线程可以用阻塞 API 发起成千上万次 web 或数据库调用而不耗尽系统资源。computation()不建议跑阻塞或 IO 密集型工作cached()不建议跑计算型工作两者在 Javadoc 里互相指向对方。执行路径一阻塞 IO 改用 cached() 或 virtual()迁移场景下最常见的任务是把subscribeOn(Schedulers.io())替换成下面两者之一。两者用法位置相同只是调度器不同。用virtual()运行阻塞调用写法来自 Whats-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 与 Whats-different-in-4.0.md 的 Cached scheduler 一节后者的前缀由rx3.io*改名为rxjava4.cached*rxjava4.cached-keep-alive-timelongcached()worker 的 keep-alive 时间默认值见CachedScheduler.KEEP_ALIVE_TIME_DEFAULT。rxjava4.cached-priorityintcached()的线程优先级默认Thread.NORM_PRIORITY。rxjava4.cached-scheduled-releaseboolean设为true时把 worker 释放模式从默认的 eager 改为 scheduled。Javadoc 说明了两模式的取舍——eager默认下 worker 立即回池、可更快复用但如果当前任务不响应中断复用可能导致延迟或死锁scheduled 下 worker 等当前任务结束才回池能减少提前复用但可能产生过多的底层 worker。rxjava4.computation-threadsintcomputation()的线程数默认为可用 CPU 数。rxjava4.computation-priorityintcomputation()的线程优先级默认Thread.NORM_PRIORITY。验证与检查看控制台输出README 示例的成功标志就是能在控制台看到 flow 的输出配合上文Thread.sleep保持 main 存活。文档没有给出固定成功日志输出内容即你的订阅回调打印的内容。阻塞保护开关RxJavaPlugins.setFailOnNonBlockingScheduler 设为true后在computation()以及single()这类非阻塞调度器上执行阻塞算子会抛出IllegalStateException。这是一个现成的手段开启它来捕获被误放到computation()上的阻塞代码。注意该设置在插件被lockdown()后不可再修改。未处理错误的去向三个调度器上未处理的错误都会交给该调度器线程的Thread.UncaughtExceptionHandlerSchedulersJavadoc 对computation()、cached()均如此说明排查挂死的流程时可以从这里入手。限制与边界cached()的 worker 线程无界 casual 使用或实现算子时必须通过Worker.dispose()释放否则可能导致系统减速或OutOfMemoryErrorJavadoc 原文。virtual()调试受限Whats-different-in-4.0.md 指出部分 IDE如 Eclipse无法正确对虚拟线程代码单步断点可能失效或挂起 IDE文档给出的替代做法是在单元测试或调试时改用传统调度器运行同一段代码。标准的Schedulers.io()已 deprecated代码里见到它时应视为待处理项而不是默认选择。应用退出前可调用Schedulers.shutdown()关闭标准调度器幂等且线程安全通过RxJavaPlugins.createXxxScheduler(ThreadFactory)创建的自定义实例则需要手动调用Scheduler.shutdown()才能让 JVM 正常退出。完整的调度器行为说明见 Scheduler.md 指向的 Scheduler 文档、Schedulers 的 Javadoc以及 Whats-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),仅供参考
返回列表