
RxJS combineLatest 深入解析Observable.prototype.combineLatest 的用法、状态机实现与测试验证【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS本文基于 RxJS v4当前仓库 package.json 中标注的版本为 4.1.0的 API 文档与源码系统讲解Rx.Observable.prototype.combineLatest(...args, [resultSelector])的调用形式、参数语义与完整示例并深入到 prototype 方法实现 与 CombineLatest 运算符内部实现结合 QUnit 单元测试 逐条验证其“最新值缓存”、完成判定与错误传播行为。读完本文你可以准确掌握 combineLatest 的发射时机与终止条件并能从源码层面理解其共享状态机state与观察者CombineLatestObserver的协作方式。API 签名与参数combineLatest的 prototype 版本挂载在Observable.prototype上签名如下Rx.Observable.prototype.combineLatest(...args, [resultSelector])参数类型说明argsarguments \| Array一组 Observable 序列可以逐个传入也可以传入一个数组[resultSelector]Function可选。每当任一源序列产生元素且所有源序列都至少发射过一次时调用省略时各源最新值组成的列表会被直接发射返回值Observable包含各源元素经resultSelector组合后的结果。其核心语义可以概括为一句话每当任一源序列产生新元素时立即用“该源的新值 其他所有源的最近一次值”调用选择器并发射结果前提是每一个源序列都至少发射过一次元素。这与zip按位置配对有本质区别——combineLatest 每次发射只“刷新”一个位置的值其余位置取缓存的最新值。prototype 版本与静态版本Rx.Observable.combineLatest的差别仅在于prototype 版本会自动把调用者自身this作为第一个源并入参数列表最终统一委托给静态实现。两个完整使用示例示例一省略 resultSelector默认发射数组两个错拍staggering的 interval 源组合不带选择器/* Have staggering intervals */ var source1 Rx.Observable.interval(100) .map(function (i) { return First: i; }); var source2 Rx.Observable.interval(150) .map(function (i) { return Second: i; }); // Combine latest of source1 and source2 whenever either gives a value with selector var source source1.combineLatest( source2 ).take(4); var subscription source.subscribe( function (x) { console.log(Next: %s, x); }, function (err) { console.log(Error: %s, err); }, function () { console.log(Completed); }); // Next: First: 0,Second: 0 // Next: First: 1,Second: 0 // Next: First: 1,Second: 1 // Next: First: 2,Second: 1 // Completed省略resultSelector时每次发射的元素是一个由所有源最新值按顺序组成的数组[s1, s2]。输出中可以看到“粘滞”特征source1 的第二个值到来时 source2 仍停留在0直到 source2 自己产生新值才刷新为1。示例二提供 resultSelector 自定义组合结果同样的两个源传入选择器函数拼接字符串/* Have staggering intervals */ var source1 Rx.Observable.interval(100) .map(function (i) { return First: i; }); var source2 Rx.Observable.interval(150) .map(function (i) { return Second: i; }); // Combine latest of source1 and source2 whenever either gives a value var source source1.combineLatest( source2, function (s1, s2) { return s1 , s2; } ).take(4); var subscription source.subscribe( function (x) { console.log(Next: %s, x); }, function (err) { console.log(Error: %s, err); }, function () { console.log(Completed); }); // Next: First: 0, Second: 0 // Next: First: 1, Second: 0 // Next: First: 1, Second: 1 // Next: First: 2, Second: 1 // Completed选择器按“源顺序”接收参数source1.combineLatest(source2, fn)中fn的第一个参数始终是 source1 的最新值第二个参数是 source2 的最新值。源码实现一prototype 方法只是参数重排combinelatestproto.js 的全部实现非常短observableProto.combineLatest function () { var len arguments.length, args new Array(len); for(var i 0; i len; i) { args[i] arguments[i]; } if (Array.isArray(args[0])) { args[0].unshift(this); // 数组形式把 this 插到数组首位 } else { args.unshift(this); // 逐参形式把 this 插到参数列表首位 } return combineLatest.apply(this, args); };它做了三件事把arguments拷贝为真正的数组args区分两种调用形式——若第一个参数是数组source1.combineLatest([obs2, obs3], fn)把thisunshift进该数组否则把thisunshift进参数列表委托给同作用域内的静态实现combineLatest定义在 combinelatest.js。由此得到两个结论prototype 版本天然支持 1 个以上源序列自身 N 个参数静态版本与 prototype 版本共用同一套CombineLatestObservable行为完全一致。源码实现二CombineLatestObservable 的共享状态机真正的运算逻辑在 combinelatest.js 中。订阅入口subscribeCore为每次订阅创建一份共享可变状态statevar state { hasValue: arrayInitialize(len, falseFactory), // 每个源是否已发射过值 hasValueAll: false, // 缓存位所有源是否都有值 isDone: arrayInitialize(len, falseFactory), // 每个源是否已完成 values: new Array(len) // 每个源的最新值缓存 }; for (var i 0; i len; i) { var source this._params[i], sad new SingleAssignmentDisposable(); subscriptions[i] sad; isPromise(source) (source observableFromPromise(source)); sad.setDisposable(source.subscribe(new CombineLatestObserver(observer, i, this._cb, state))); } return new NAryDisposable(subscriptions);这里有几个值得注意的实现细节每个源都分配一个独立的CombineLatestObserver但它们共享同一个state对象——“最新值缓存”就体现在values数组上各观察者在自己的next中只更新自己下标i对应的位置Promise 源被自动包装isPromise(source) (source observableFromPromise(source))因此 combineLatest 可以混合 Observable 与 Promise这与 TypeScript 签名中的ObservableOrPromiseT对应见下文hasValueAll是一个缓存标志位一旦所有源都有值便置为true后不再回查避免每次next都执行hasValue.every(identity)所有源订阅被收集进一个NAryDisposable任一订阅失败或外部取消时都会统一处置避免悬挂订阅。combineLatest静态入口同时负责解析resultSelector最后一个参数是函数则弹出作为选择器否则回退为argumentsToArray把各源最新值组装成数组发射——这正对应文档中“If omitted, a list with the elements will be yielded”的行为。源码实现三CombineLatestObserver 的发射、完成与错误规则CombineLatestObserver是理解 combineLatest 行为的钥匙CombineLatestObserver.prototype.next function (x) { this._state.values[this._i] x; this._state.hasValue[this._i] true; if (this._state.hasValueAll || (this._state.hasValueAll this._state.hasValue.every(identity))) { var res tryCatch(this._cb).apply(null, this._state.values); if (res errorObj) { return this._o.onError(res.e); } this._o.onNext(res); } else if (this._state.isDone.filter(notTheSame(this._i)).every(identity)) { self._o.onCompleted(); } }; CombineLatestObserver.prototype.error function (e) { this._o.onError(e); }; CombineLatestObserver.prototype.completed function () { this._state.isDone[this._i] true; this._state.isDone.every(identity) this._o.onCompleted(); };注原文代码中onCompleted分支为this._o.onCompleted();此处self为笔误实际源码见 combinelatest.js#L56-L65。从中可以提炼出三条规则发射门槛只有当hasValue.every为真即所有源至少发射过一次时next才会调用选择器并onNext。任何“迟到”之前早到的源的值只会被静默缓存进values不会触发发射提前完成判定若当前源尚未产生过值门槛未满足但除自己之外的所有源都已完成isDone.filter(notTheSame(this._i)).every(identity)则直接onCompleted——因为其他源都不会再有新值门槛永远无法满足序列提前终止错误与选择器异常任一源onError会立即透传给下游error方法直传选择器函数自身抛出的异常由tryCatch捕获后转为onError不会让订阅崩溃。completed中则要求所有源都完成才触发onCompleted。测试用例验证完成与错误的边界行为tests/observable/combinelatest.js 使用TestScheduler 热hotObservable 对边界行为做了系统验证。热序列的时间线里t150的发射发生在默认订阅时刻t200之前因此只有t200后的事件真正参与运算。下面挑几个关键用例解读。“interleaved with tail”最新值缓存的完整轨迹这是最能体现“latest”语义的用例tests/observable/combinelatest.js#L454-L481e1onNext(150, 1)、onNext(215, 2)、onNext(225, 4)、onCompleted(230)e2onNext(150, 1)、onNext(220, 3)、onNext(230, 5)、onNext(235, 6)、onCompleted(250)订阅发生在t200t150的两个值均未被观察。随后时间事件state 变化输出215e1 发射 2values[2, ?]hasValueAll 仍为 falsee2 无值仅缓存无220e2 发射 3values[2, 3]hasValueAll 置 true首次满足门槛onNext(220, 23)225e1 发射 4values[4, 3]onNext(225, 43)230e1 完成、e2 发射 5isDone[0]true未全部完成不终止values[4, 5]onNext(230, 45)235e2 发射 6values[4, 6]onNext(235, 46)240e2 发射 7values[4, 7]onNext(240, 47)250e2 完成isDone 全 trueonCompleted(250)注意onNext(230, 45)e1 与 e2 在同一时刻230分别触发“完成”和“发射”e1 虽已完成其最新值 4 依然被用于后续组合——已完成源的最终值会被保留到最后。never / empty / return 组合矩阵测试覆盖了 never永不发射、empty立即完成、return发射后完成三者的全部两两组合结论与源码规则一一对应never never、never empty、empty never、never return、return never门槛永远无法满足 →零输出tests/observable/combinelatest.js#L15-L143empty empty两源在 210 完成isDone全 true →onCompleted(210)empty returne1empty210 完成e2 在 215 发射 2但 e1 从未产生值且其余源e1已完成命中“提前完成判定” →onCompleted(215)215 的值被丢弃。错误传播规则首个错误胜出throw throw用例中 e1 在 220 出错、e2 在 230 出错最终只发出onError(220, error1)错误不受“完成”阻挡throw after complete left中 e1 已于 220 完成e2 在 230 出错依然输出onError(230, error)tests/observable/combinelatest.js#L412-L452选择器抛错selector throws用例验证选择器内部throw error会经tryCatch转为onError(220, error)tests/observable/combinelatest.js#L557-L577无选择器时发射数组return return no selector用例中输出为onNext(220, [2, 3])即默认选择器argumentsToArray的行为tests/observable/combinelatest.js#L166-L185。TypeScript 类型签名最多 9 个源的完整重载在 ts/core/linq/observable/combinelatestproto.ts 中combineLatest提供了从 1 个到 9 个源序列的重载以及数组形式重载类型层面把“源顺序即参数顺序”的约定固化下来combineLatestT2, T3, T4, T5, T6, T7, T8, T9( second: ObservableOrPromiseT2, third: ObservableOrPromiseT3, /* ... 其余源参数 ... */ ninth: ObservableOrPromiseT9 ): Observable[T, T2, T3, T4, T5, T6, T7, T8, T9]; combineLatestTOther(sources: ObservableOrPromiseTOther[]): ObservableTOther[];类型定义文件末尾还内置了编译期验证示例明确了四种典型用法ts/core/linq/observable/combinelatestproto.ts#L188-L201var r1: Rx.Observable{ vo: boolean, vio: string, vp: { a: string }, vso: number } o.combineLatest(io, p, so, (vo, vio, vp, vso) ({ vo, vio, vp, vso })); var r2: Rx.Observable[boolean, string, { a: string }, number] o.combineLatest(io, p, so); var r3: Rx.Observablenumber o.combineLatest(any[][io, so, p], (v1, items) 5); var r4: Rx.Observableany[] o.combineLatest(any[][io, so, p]);从这些类型声明可以看出每个参数类型都是ObservableOrPromiseT即Promise 可以直接作为源传入与源码中的observableFromPromise包装逻辑对应带选择器时返回值类型由选择器决定r1不带选择器时返回元组类型r2静态版本Rx.Observable.combineLatest在 ts/core/linq/observable/combinelatest.ts 中有对应的同构重载其末尾的编译期示例同样演示了选择器、元组、数组选择器、数组无选择器四种形态。仓库中的相关实现与测试位置内容路径prototype 方法参数重排并委托静态实现src/core/linq/observable/combinelatestproto.js静态实现CombineLatestObservableCombineLatestObserversrc/core/perf/operators/combinelatest.js模块化CommonJS实现供按需 requiresrc/modular/observable/combinelatest.jsQUnit 单元测试TestScheduler 时间线用例tests/observable/combinelatest.jsprototype 版本 TypeScript 重载ts/core/linq/observable/combinelatestproto.ts静态版本 TypeScript 重载ts/core/linq/observable/combinelatest.tsNPM 包定义rxv4.1.0package.jsonAPI 文档原文doc/api/core/operators/combinelatestproto.md模块化实现src/modular/observable/combinelatest.js与核心实现逻辑一致同样包含hasValue/hasValueAll/isDone/values四元状态、Promise 包装与tryCatch保护仅将依赖改为 CommonJSrequire形式可在 src/modular/index.js 中按需组合。小结combineLatest是 RxJS v4 中用于合并多个“持续变化的状态源”的核心运算符理解它的关键在于三点发射门槛所有源至少各发射一次后才开始输出早到的值只进缓存hasValue/values粘滞组合此后任一源的新值都会触发一次“新值 其余源最新值”的组合发射已完成源的最终值会被保留使用终止语义任一源出错立即传播所有源完成后正常完成若某源尚未产生值而其余源已全部完成则提前onCompleted并丢弃该值。这些行为既可通过 API 文档 的两个 interval 示例直观感受也可在 tests/observable/combinelatest.js 的时间线用例与 源码状态机 中找到逐条对应的证据。【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考