恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
RxJava 组合流算子详解:startWith、merge、zip、combineLatest 与 switchOnNext 的完整实战指南
首页
资讯中心
/
RxJava 组合流算子详解:startWith、merge、zip、combineLatest 与 switchOnNext 的完整实战指南
RxJava 组合流算子详解:startWith、merge、zip、combineLatest 与 switchOnNext 的完整实战指南
发布时间:2026/9/18 20:57:21
RxJava 组合流算子详解startWith、merge、zip、combineLatest 与 switchOnNext 的完整实战指南【免费下载链接】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 仓库中的 docs/Combining-Observables.md 展开系统讲解用于组合多个流Observable/Flowable 等的 8 类算子startWith、merge、mergeDelayError、zip、combineLatest、switchOnNext、join/groupJoin以及可选包rxjava-joins提供的and/then/when。读完本文你将掌握每个算子的语义边界、在Flowable/Observable/Maybe/Single/Completable五种源类型上的可用性并理解这些算子在当前仓库源码io.reactivex.rxjava4包中的真实签名与错误传播行为。组合算子总览RxJava 组合类算子解决的问题是多个独立的事件流如何变成一个。仓库文档列出的 8 个条目如下继承自原文档目录算子语义一句话所属startWith在源流开始发射前先发射一组预置值标准算子merge合并多个流任一源的onError立即终止整个合并流标准算子mergeDelayError合并多个流但onError被延迟到所有源完成后才下发标准算子zip按下标将两个或多个流的第 n 个值配对经函数处理后发射标准算子combineLatest任一源发射时取各源最新值组合后经函数发射标准算子switchOnNext将“发射流的流”折叠为单个流只保留最近一个内层流的输出标准算子join/groupJoin当一个源的项落入另一个源项所给出的时间窗口内时将两者关联标准算子and/then/when通过Pattern/Plan中间对象组合多个流的值集合可选rxjava-joins包五种源类型上的可用性矩阵原文档为每个算子标注了它支持哪类源checkmark 标记。整理成表格如下算子FlowableObservableMaybeSingleCompletablestartWith✓✓———merge✓✓✓✓✓mergeDelayError✓✓✓✓✓zip✓✓✓✓—combineLatest✓✓———switchOnNext✓✓———注意两个规律Maybe/Single/Completable都是“最多 1 个值/0 个值”的源因此凡是依赖“流式多值”语义的算子startWith、combineLatest、switchOnNext对它们没有意义Completable连值都没有所以zip这类“按值组合”的算子也不适用。merge和mergeDelayError则对五种源全量支持。startWith先发射预置值再接源流startWith的语义是在开始发射源 Observable 的项之前先发射一组指定的项。原文档示例ObservableString names Observable.just(Spock, McCoy); names.startWith(Kirk).subscribe(item - System.out.println(item)); // prints Kirk, Spock, McCoy从当前仓库源码看Observable.java 中startWith家族远比“插一个值”丰富共提供了 7 个相关入口startWithIterable先发射一个Iterable中的全部项再发射源流startWith(CompletableSource)先挂接一个Completable它会以onComplete结束等价于“什么都不先发射仅保证顺序”再发射源流startWith(SingleSource) 与 startWith(MaybeSource)先发射Single/Maybe的结果值若Maybe为空则跳过再发射源流startWith(ObservableSource)先完整发射另一个Observable再发射源流startWithItem 与 startWithArray先发射单个值或一组值原文档示例names.startWith(Kirk)对应的就是这类入口。这意味着startWith在 RxJava 4 中是“把任意单值型/多值型源前置拼接”的统一入口而不仅仅是插常量——把另一个ObservableSource当作前置流时其onError会直接传递并终止结果流行为与concat的前半段一致。对应测试可参考 ObservableStartWithTests.java。merge无保序合并错误立即传播merge将多个 Observable 合并为一个不保证各源之间的相对顺序。原文档示例使用mergeWithObservable.just(1, 2, 3) .mergeWith(Observable.just(4, 5, 6)) .subscribe(item - System.out.println(item)); // prints 1, 2, 3, 4, 5, 6错误语义是merge的关键设计点任何一个源发出的onError都会立即透传给下游观察者并终止整个合并后的 Observable。源码层面Observable.java 提供了多组静态/实例入口可按输入形态选择merge(Iterable) 及其带StandardConcurrentBufferedConfig的并发变体输入是“源的集合”merge(ObservableSource? extends ObservableSource) 及并发变体输入是“发射源”的高阶流类似flatMap的同步版本但保留内层订阅关系mergeArray(ObservableSource...) 及其配置变体输入是编译期已知的参数数组适合 2~3 个固定流合并实例方法 mergeWith(ObservableSource)把自身与另一个ObservableSource合并内部实现就是mergeArray(this, other)。值得注意的演进当前仓库的mergeWith除了接受ObservableSource还重载了 mergeWith(SingleSource)、mergeWith(MaybeSource) 和 mergeWith(CompletableSource)即可以把一个单值型源当作“零头”并入流中。由于merge不占用特定 SchedulerJavadoc 明确标注“does not operate by default on a particular Scheduler”它不会引入线程切换是成本最低的合并手段。相关测试见 ObservableMergeTests.java。mergeDelayError错误延迟到全部源完成mergeDelayError与merge的组合逻辑相同但错误处理策略不同任何源发出的onError都被扣留直到所有被合并的 Observable 全部完成才下发给下游观察者。原文档示例ObservableString observable1 Observable.error(new IllegalArgumentException()); ObservableString observable2 Observable.just(Four, Five, Six); Observable.mergeDelayError(observable1, observable2) .subscribe(item - System.out.println(item)); // emits Four, Five, Six然后是那个 IllegalArgumentException // 在本例中由于 subscribe 只写了 onNext 未处理错误会抛出 OnErrorNotImplementedException这个行为非常适合“多个相互独立的请求宁可拿到部分结果、最后再统一处理失败”的场景。从源码 Javadoc 可以确认更细致的边界以 Maybe.java 中的说明为例即使多个被合并的源都发送了onErrormergeDelayError也只会向下游传递第一个错误Reactive Streams 规范本身也只允许一次onError。另外Maybe/Single各自至多 0/1 个值天然不存在“高阶源”问题因此仓库 Javadoc 特别注明“不需要mergeDelayError(MaybeSourceMaybeSourceT)这种重载”——这解释了可用性矩阵中Maybe/Single/Completable均被勾选的原因。仓库中对应的内部实现类位于 OperatorMergeArray.java 所在包io.reactivex.rxjava4.internal.operators.observableMaybe/Single/Completable的mergeDelayErrorJavadoc 均带有各自的算子时序图注释如 Single.java。zip按下标对齐组合zip将两个或多个 Observable 发出的项集合通过指定函数两两组合并按组合结果发射。它的核心是位置对齐第 n 个输出由第 n 个输入值决定任何一边多出的尾部值会被丢弃以较短的流为准。原文档示例ObservableString firstNames Observable.just(James, Jean-Luc, Benjamin); ObservableString lastNames Observable.just(Kirk, Picard, Sisko); firstNames.zipWith(lastNames, (first, last) - first last) .subscribe(item - System.out.println(item)); // prints James Kirk, Jean-Luc Picard, Benjamin Sisko与combineLatest的关键区别zip是“第 n 对第 n”combineLatest是“谁动谁触发、取各自最新值”。zip适用于两边数量天然一致的数据如首名/姓氏而combineLatest适用于两边节奏不同、但都关心最新状态的场景。在可用性上zip支持Flowable/Observable/Maybe/Single四类源Completable无值可 zip。测试参考 ObservableZipTests.java。combineLatest最新值组合当两个 Observable 中任意一个发射项时combineLatest取每个源最近一次发射的值经指定函数组合后发射。原文档示例用两个不同节拍的interval演示了“快流带动慢流”的行为ObservableLong newsRefreshes Observable.interval(100, TimeUnit.MILLISECONDS); ObservableLong weatherRefreshes Observable.interval(50, TimeUnit.MILLISECONDS); Observable.combineLatest(newsRefreshes, weatherRefreshes, (newsRefreshTimes, weatherRefreshTimes) - Refreshed news newsRefreshTimes times and weather weatherRefreshTimes) .subscribe(item - System.out.println(item)); // prints: // Refreshed news 0 times and weather 0 // Refreshed news 0 times and weather 1 // Refreshed news 0 times and weather 2 // Refreshed news 1 times and weather 2 // Refreshed news 1 times and weather 3 // ...从输出可以读出三条行为规则(1) 某源尚未发射过任何值时组合不会发生——所以第一次输出必须等两边都有值(2) 快流weather50ms每次发射都会触发一次组合而慢流news100ms的值保持“最近一次”不变(3) 各源完成后不终止结果直到所有源都onComplete才onComplete。实际开发中常用技巧是若担心“慢流一直没值导致不输出”可以用一个快速流如interval或timer作为触发器。combineLatest仅适用于Flowable/Observable矩阵中Maybe/Single/Completable均不支持因为单值源“最新值”语义退化。测试参考 ObservableCombineLatestTests.java。switchOnNext只保留最新的内层流switchOnNext把一个发射 Observable 的 Observable即ObservableObservableT折叠为单一 Observable只发射最近一次发出的内层 Observable的项此前的内层流被取消订阅。原文档示例外层每 1s 产生一个新内层 intervalObservableObservableString timeIntervals Observable.interval(1, TimeUnit.SECONDS) .map(ticks - Observable.interval(100, TimeUnit.MILLISECONDS) .map(innerInterval - outer: ticks - inner: innerInterval)); Observable.switchOnNext(timeIntervals) .subscribe(item - System.out.println(item)); // prints: // outer: 0 - inner: 0 // outer: 0 - inner: 1 // ... // outer: 0 - inner: 8 // outer: 1 - inner: 0 // 1 秒后切换到新的内层流 // outer: 1 - inner: 1 // ...这是实现“用户快速输入时取消上一次请求”搜索建议、防抖请求的经典算子外层每次发射新请求流旧请求流即刻被丢弃。与merge/concat对比merge会同时运行所有内层流concat按序排队只有switchOnNext是“抢占式”——最新者胜。它同样只在Flowable/Observable上可用Maybe/Single不能成为内层流。join / groupJoin 与 rxjava-joins 包原文档最后两节描述了两类“表连接式”的组合算子此处如实继承其定义join()与groupJoin()当来自一个 Observable 的项落入由另一个 Observable 项所指定的时间窗口内时把两边的项关联起来输出groupJoin输出的是该窗口内项的集合。and()、then()、when()通过Pattern和Plan两个中间对象组合两个或多个 Observable 发出的项集合。需要注意rxjava-joins标注的含义原文档原文and/then/when属于可选的rxjava-joins包位于rxjava-contrib之下不包含在标准 RxJava 算子集合中。当前仓库的src/main下未包含该模块因此这三个算子需要额外引入对应依赖才能使用本文仅记录其语义与归属不展开其源码。在五种源类型上如何选择结合前文的可用性矩阵实践中可以按源类型快速决策Flowable/Observable多值流8 类组合能力基本全开按“是否保序、是否对齐、是否抢占”选merge、zip、switchOnNext错误策略二选一merge立即失败或mergeDelayError拿满结果再报错。Maybe/Single至多 1 值可用merge/mergeDelayError/zip“合并”在这里更接近“并行发起、按完成聚合”mergeDelayError的价值尤其突出——多个Single请求并发执行全部返回后统一交付错误。Completable0 值只有merge/mergeDelayError用于并发执行多个“纯完成信号”任务如多个清理/初始化动作。小结startWith在当前仓库中演化出 7 个入口Observable.java可前置Iterable、数组、单项乃至任意ObservableSource/MaybeSource/SingleSourcemerge立即传播onError且不切线程mergeDelayError扣留错误、全部源完成后只下发第一个错误zip按下标对齐、combineLatest按最新值组合二者是“对齐”与“触发”两种不同的心智模型switchOnNext提供抢占式内层流切换join/groupJoin是时间窗口连接and/then/when属于可选rxjava-joins包。进一步阅读可参考仓库中的 docs/Operator-Matrix.md完整算子矩阵、docs/Async-Operators.md异步组合算子以及src/test/java/io/reactivex/rxjava4/observable/与src/test/java/io/reactivex/rxjava4/flowable/下对应的ObservableMergeTests、ObservableZipTests、ObservableCombineLatestTests、ObservableStartWithTests、ObservableSwitchOnNext相关测试用测试用例验证上述各算子的边界行为。【免费下载链接】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),仅供参考