RxJS combineLatest 深入解析:Observable.prototype.combineLatest 的用法、状态机实现与测试验证
2026/9/20 10:14:18 网站建设 项目流程

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); };

它做了三件事:

  1. arguments拷贝为真正的数组args
  2. 区分两种调用形式——若第一个参数是数组(source1.combineLatest([obs2, obs3], fn)),把thisunshift进该数组;否则把thisunshift进参数列表;
  3. 委托给同作用域内的静态实现combineLatest(定义在 combinelatest.js)。

由此得到两个结论:prototype 版本天然支持 1 个以上源序列(自身 + N 个参数);静态版本与 prototype 版本共用同一套CombineLatestObservable,行为完全一致。

源码实现二:CombineLatestObservable 的共享状态机

真正的运算逻辑在 combinelatest.js 中。订阅入口subscribeCore为每次订阅创建一份共享可变状态state

var 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 签名中的ObservableOrPromise<T>对应,见下文);
  • 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。)

从中可以提炼出三条规则:

  1. 发射门槛:只有当hasValue.every为真(即所有源至少发射过一次)时,next才会调用选择器并onNext。任何“迟到”之前早到的源的值只会被静默缓存进values,不会触发发射;
  2. 提前完成判定:若当前源尚未产生过值(门槛未满足),但除自己之外的所有源都已完成isDone.filter(notTheSame(this._i)).every(identity)),则直接onCompleted——因为其他源都不会再有新值,门槛永远无法满足,序列提前终止;
  3. 错误与选择器异常:任一源onError会立即透传给下游(error方法直传);选择器函数自身抛出的异常由tryCatch捕获后转为onError,不会让订阅崩溃。completed中则要求所有源都完成才触发onCompleted

测试用例验证:完成与错误的边界行为

tests/observable/combinelatest.js 使用TestScheduler+ 热(hot)Observable 对边界行为做了系统验证。热序列的时间线里t=150的发射发生在默认订阅时刻t=200之前,因此只有t>=200后的事件真正参与运算。下面挑几个关键用例解读。

“interleaved with tail”:最新值缓存的完整轨迹

这是最能体现“latest”语义的用例(tests/observable/combinelatest.js#L454-L481):

  • e1:onNext(150, 1)onNext(215, 2)onNext(225, 4)onCompleted(230)
  • e2:onNext(150, 1)onNext(220, 3)onNext(230, 5)onNext(235, 6)onCompleted(250)

订阅发生在t=200t=150的两个值均未被观察。随后:

时间事件state 变化输出
215e1 发射 2values=[2, ?],hasValueAll 仍为 false(e2 无值),仅缓存
220e2 发射 3values=[2, 3],hasValueAll 置 true,首次满足门槛onNext(220, 2+3)
225e1 发射 4values=[4, 3]onNext(225, 4+3)
230e1 完成、e2 发射 5isDone[0]=true(未全部完成,不终止);values=[4, 5]onNext(230, 4+5)
235e2 发射 6values=[4, 6]onNext(235, 4+6)
240e2 发射 7values=[4, 7]onNext(240, 4+7)
250e2 完成isDone 全 trueonCompleted(250)

注意onNext(230, 4+5):e1 与 e2 在同一时刻(230)分别触发“完成”和“发射”,e1 虽已完成,其最新值 4 依然被用于后续组合——已完成源的最终值会被保留到最后

never / empty / return 组合矩阵

测试覆盖了 never(永不发射)、empty(立即完成)、return(发射后完成)三者的全部两两组合,结论与源码规则一一对应:

  • never + nevernever + emptyempty + nevernever + returnreturn + never:门槛永远无法满足 →零输出(tests/observable/combinelatest.js#L15-L143);
  • empty + empty:两源在 210 完成,isDone全 true →onCompleted(210)
  • empty + return:e1(empty)210 完成,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 个源序列的重载,以及数组形式重载,类型层面把“源顺序即参数顺序”的约定固化下来:

combineLatest<T2, T3, T4, T5, T6, T7, T8, T9>( second: ObservableOrPromise<T2>, third: ObservableOrPromise<T3>, /* ... 其余源参数 ... */ ninth: ObservableOrPromise<T9> ): Observable<[T, T2, T3, T4, T5, T6, T7, T8, T9]>; combineLatest<TOther>(sources: ObservableOrPromise<TOther>[]): Observable<TOther[]>;

类型定义文件末尾还内置了编译期验证示例,明确了四种典型用法(ts/core/linq/observable/combinelatestproto.ts#L188-L201):

var 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.Observable<number> = o.combineLatest(<any[]>[io, so, p], (v1, items) => 5); var r4: Rx.Observable<any[]> = o.combineLatest(<any[]>[io, so, p]);

从这些类型声明可以看出:

  • 每个参数类型都是ObservableOrPromise<T>,即Promise 可以直接作为源传入,与源码中的observableFromPromise包装逻辑对应;
  • 带选择器时返回值类型由选择器决定(r1),不带选择器时返回元组类型(r2);
  • 静态版本Rx.Observable.combineLatest在 ts/core/linq/observable/combinelatest.ts 中有对应的同构重载,其末尾的编译期示例同样演示了选择器、元组、数组选择器、数组无选择器四种形态。

仓库中的相关实现与测试位置

内容路径
prototype 方法(参数重排并委托静态实现)src/core/linq/observable/combinelatestproto.js
静态实现:CombineLatestObservable+CombineLatestObserversrc/core/perf/operators/combinelatest.js
模块化(CommonJS)实现,供按需 requiresrc/modular/observable/combinelatest.js
QUnit 单元测试(TestScheduler 时间线用例)tests/observable/combinelatest.js
prototype 版本 TypeScript 重载ts/core/linq/observable/combinelatestproto.ts
静态版本 TypeScript 重载ts/core/linq/observable/combinelatest.ts
NPM 包定义(rx,v4.1.0)package.json
API 文档原文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 中用于合并多个“持续变化的状态源”的核心运算符,理解它的关键在于三点:

  1. 发射门槛:所有源至少各发射一次后才开始输出,早到的值只进缓存(hasValue/values);
  2. 粘滞组合:此后任一源的新值都会触发一次“新值 + 其余源最新值”的组合发射,已完成源的最终值会被保留使用;
  3. 终止语义:任一源出错立即传播;所有源完成后正常完成;若某源尚未产生值而其余源已全部完成,则提前onCompleted并丢弃该值。

这些行为既可通过 API 文档 的两个 interval 示例直观感受,也可在 tests/observable/combinelatest.js 的时间线用例与 源码状态机 中找到逐条对应的证据。

【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

立即咨询