如何用 RxJava 的 ParallelFlowable 做并行数据处理?parallel、runOn 与 sequential 用法及适用边界 如何用 RxJava 的 ParallelFlowable 做并行数据处理parallel、runOn 与 sequential 用法及适用边界【免费下载链接】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 的响应式风格时ParallelFlowable是 RxJava 提供的并行处理入口。它把一个Flowable源拆成多条并行的 rail轨道在多条线程上并行执行map、filter、flatMap等操作最后再合并回一个普通的Flowable。本文基于项目文档 Parallel-flows.md 和源码 ParallelFlowable.java、Flowable.java 的 Javadoc给出可执行的用法和明确的适用边界。先明确前提避免走错方向ParallelFlowable自 RxJava 2.0.5 引入experimental2.1 转 beta。它不是一个新的响应式基类而是Flowable的并行模式子领域特定语言。因此不存在ParallelObservable——官方文档的解释是并行化的目的通常是单线程处理数据太慢而避免内部队列被压垮必须依赖背压Observable没有背压所以并行模式只在支持背压的Flowable上提供。依赖以 Gradle 方式引入见 README.md 的 Getting started 一节implementation io.reactivex.rxjava4:rxjava:4.x.y其中4.x.y需按仓库 README 的说明替换为 Maven Central 上的最新版本号。三步主路径parallel → runOn → sequential并行处理的完整链路是固定的三步文档给出的示例如下// 1. 从顺序源进入并行世界 ParallelFlowableInteger source Flowable.range(1, 1000).parallel(); // 2. 指定每条 rail 在哪个 Scheduler 上执行引入异步消费 ParallelFlowableInteger psource source.runOn(Schedulers.io()); // 3. 应用并行操作符后合并回普通 Flowable FlowableInteger result psource.filter(v - v % 3 0).map(v - v * v).sequential();第一步Flowable.parallel()拆轨parallel()创建多条 rail并把上游值以 round-robin轮询方式分派给各 rail。源码中有三个重载见 Flowable.javaparallel() // 默认rail 数 Runtime.getRuntime().availableProcessors() parallel(int parallelism) // 指定 rail 数 parallel(int parallelism, int prefetch) // 指定 rail 数与每条 rail 的预取量默认并行度是可用 CPU 数默认从顺序源预取Flowable.bufferSize()128个值两者都可以通过上面的重载指定。parallelism和prefetch传非正数会抛IllegalArgumentException。parallel()本身不会引入异步执行——它只是准备好并行流。rail 不会自己并行起来必须接着调用runOn(Scheduler)指定每条 rail 跑在哪个线程上。这是源码 Javadoc 中反复强调的一点Flowable.parallel()的注释the rails dont execute in parallel on their own and one needs to applyrunOn(Scheduler)。除Flowable上调用外也可以直接对任意Publisher用静态工厂ParallelFlowable.from(publisher)及其(parallelism)、(parallelism, prefetch)重载见 ParallelFlowable.java。第二步runOn(Scheduler)引入并行执行runOn指定每条 rail 在哪里观察并处理值。关键行为见 ParallelFlowable.java Javadoc会按并行度调用Scheduler.createWorker()次数等于 rail 数即 rail 数与Scheduler自身的并行度不必一致。如果Scheduler的并行度低于ParallelFlowable的并行度部分 rail 可能落在同一个线程/worker 上——此时并不会报错只是实际线程数不足。该操作符自带内置 trampoline 逻辑不要求Scheduler是 trampolining 的。runOn(scheduler, prefetch)重载可指定每条 rail 的预取量默认同样是Flowable.bufferSize()。文档给出的三类典型选型任务特征SchedulerCPU 密集型计算Schedulers.computation()阻塞 / IO 型任务Schedulers.io()单元测试TestScheduler第三步sequential()合并回 Flowable并行操作完成后用sequential()把各 rail 合并回单个Flowable。相关变体见 ParallelFlowable.javasequential() // 默认预取量合并 sequential(int prefetch) // 指定每条 rail 的预取量 sequentialDelayError() // 延迟所有 rail 的 error直到全部终止 sequentialDelayError(int prefetch) sorted(Comparator? super T comparator) // 合并成按比较器排好序的顺序流两个必须知道的点sequential()不保证任何顺序值经过并行操作符后以 round-robin 或同序方式合并最终序列的顺序不做承诺。如果你的下游逻辑依赖顺序不要依赖sequential()应改用sorted(comparator)要求源是有限的ParallelFlowable。sorted()和toSortedList()都要求源是有限的requires a finite source对无限流不能使用。并行模式支持哪些操作符文档列举了并行模式下可用的少数选定操作符map、filter、concatMap、flatMap、collect、reduce等。从 ParallelFlowable.java 完整签名还能确认flatMapIterable、flatMapStream与concatMapStream等价、mapOptional、reduce、collect(Collector)、fromArray(Publisher...)等。filter和doOnNext提供带ParallelFailureHandling或BiFunctionLong, Throwable, ParallelFailureHandling错误处理函数的重载用于决定 rail 上抛出异常后是继续还是终止。使用map/filter/reduce时的硬性约束写在 Javadoc 里同一个函数可能被多个线程并发调用the same mapper/reducer function may be called from multiple threads concurrently你的映射函数、谓词、归约函数必须线程安全或无共享可变状态。完整示例与结果验证仓库测试代码 ParallelFlowableTest.java 给出了标准写法主路径等价于FlowableInteger source Flowable.range(1, 100).hide(); ExecutorService exec Executors.newFixedThreadPool(4); Scheduler scheduler Schedulers.from(exec); FlowableInteger result ParallelFlowable.from(source, 4) // 4 条 rail .runOn(scheduler) .map(v - v 1) .sequential(); TestSubscriberInteger ts new TestSubscriber(); result.subscribe(ts);验证方式沿用测试代码中的TestSubscriberEx断言模式assertSubscribed().assertValueCount(n).assertComplete().assertNoErrors()测试中并行模式还会先awaitDone(10, TimeUnit.SECONDS)等待完成。由于sequential()不保证顺序对map(v - v 1)这类示例正确性判定是收到恰好 100 个值、无 error、正常onComplete且值的集合等于{2..101}——而不是逐位比较顺序。示例中exec用完需要shutdown()测试代码在finally块中执行。如果不调用runOn直接sequential()即文档中的 sequential mode值仍在订阅线程上处理测试 ParallelFlowableTest.java 对 1 到 32 的并行度都验证了 100 万条值全部送达且无 error。适用边界缺少常见操作符take、skip等许多常规操作符在ParallelFlowable上不可用文档原文several typical operators such astake,skipand many others are not available。需要take语义时只能在sequential()之后于普通Flowable上使用。不是新基类它是Flowable的并行 DSL没有ParallelObservable并行路径全程依赖背压防止内部队列被压垮。线程模型rail 数不需要等于Scheduler的线程数但Scheduler线程更少时部分 rail 会共用线程runOn的 worker 数严格等于并行度每条 rail 一个 worker。参数校验parallel()/from()的parallelism、prefetchrunOn的prefetch非正数都会抛IllegalArgumentException。并发回调map/filter/reduce等函数的同实例会被多线程并发调用共享可变状态的函数会破坏正确性。直接订阅不调用sequential()而直接对ParallelFlowable调用subscribe(Subscriber[] subscribers)时订阅者数组长度必须等于parallelism()否则每个订阅者会收到IllegalArgumentExceptionvalidate逻辑见 ParallelFlowable.java。错误传播默认sequential()/ 普通flatMap等错误立即终止需要全部跑完再报错误时改用sequentialDelayError()、concatMapDelayError(mapper, tillTheEnd)或flatMap(mapper, delayError, ...)。延伸并行操作符的完整行为细节以源码 Javadoc 为准ParallelFlowable.java。各操作符的测试覆盖了顺序/并行两种模式、多种 rail 数与TestScheduler可作为写法参照ParallelFlowableTest.java同目录下还有 ParallelRunOnTest.java、ParallelCollectTest.java 等。Parallel-flows.md 中 Parallel operators 一节目前标注为 TBD操作符级细节需以上述 Javadoc 与测试为准。【免费下载链接】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),仅供参考