如何用 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 引入(experimental,2.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. 从顺序源进入并行世界 ParallelFlowable<Integer> source = Flowable.range(1, 1000).parallel(); // 2. 指定每条 rail 在哪个 Scheduler 上执行(引入异步消费) ParallelFlowable<Integer> psource = source.runOn(Schedulers.io()); // 3. 应用并行操作符后,合并回普通 Flowable Flowable<Integer> result = psource.filter(v -> v % 3 == 0).map(v -> v * v).sequential();第一步:Flowable.parallel()拆轨
parallel()创建多条 rail,并把上游值以 round-robin(轮询)方式分派给各 rail。源码中有三个重载(见 Flowable.java):
parallel() // 默认: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 don't 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())。
文档给出的三类典型选型:
| 任务特征 | Scheduler |
|---|---|
| CPU 密集型计算 | Schedulers.computation() |
| 阻塞 / IO 型任务 | Schedulers.io() |
| 单元测试 | TestScheduler |
第三步:sequential()合并回 Flowable
并行操作完成后,用sequential()把各 rail 合并回单个Flowable。相关变体(见 ParallelFlowable.java):
sequential() // 默认预取量合并 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或BiFunction<Long, Throwable, ParallelFailureHandling>错误处理函数的重载,用于决定 rail 上抛出异常后是继续还是终止。
使用map/filter/reduce时的硬性约束写在 Javadoc 里:同一个函数可能被多个线程并发调用(the same mapper/reducer function may be called from multiple threads concurrently),你的映射函数、谓词、归约函数必须线程安全或无共享可变状态。
完整示例与结果验证
仓库测试代码 ParallelFlowableTest.java 给出了标准写法,主路径等价于:
Flowable<Integer> source = Flowable.range(1, 100).hide(); ExecutorService exec = Executors.newFixedThreadPool(4); Scheduler scheduler = Schedulers.from(exec); Flowable<Integer> result = ParallelFlowable.from(source, 4) // 4 条 rail .runOn(scheduler) .map(v -> v + 1) .sequential(); TestSubscriber<Integer> 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、prefetch,runOn的prefetch非正数都会抛IllegalArgumentException。 - 并发回调:
map/filter/reduce等函数的同实例会被多线程并发调用,共享可变状态的函数会破坏正确性。 - 直接订阅:不调用
sequential()而直接对ParallelFlowable调用subscribe(Subscriber[] subscribers)时,订阅者数组长度必须等于parallelism(),否则每个订阅者会收到IllegalArgumentException(validate逻辑见 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),仅供参考