如何用 RxJava 的 ParallelFlowable 做并行数据处理?parallel、runOn 与 sequential 用法及适用边界
2026/9/10 3:14:02 网站建设 项目流程

如何用 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"(轨道),在多条线程上并行执行mapfilterflatMap等操作,最后再合并回一个普通的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)个值,两者都可以通过上面的重载指定。
  • parallelismprefetch传非正数会抛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),对无限流不能使用。

并行模式支持哪些操作符

文档列举了并行模式下可用的"少数选定操作符":mapfilterconcatMapflatMapcollectreduce等。从 ParallelFlowable.java 完整签名还能确认flatMapIterableflatMapStream(与concatMapStream等价)、mapOptionalreducecollect(Collector)fromArray(Publisher...)等。filterdoOnNext提供带ParallelFailureHandlingBiFunction<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。

适用边界

  • 缺少常见操作符takeskip等许多常规操作符在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()parallelismprefetchrunOnprefetch非正数都会抛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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询