RxJS 4 中的Rx.Observable.for:基于数组批量生成并串联观察序列的完整指南
【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS
导读
Rx.Observable.for(sources, resultSelector, [thisArg])是 RxJS(Reactive Extensions for JavaScript,本仓库对应 RxJS v4 系列)中一个静态工厂操作符:它把普通数组中的每一个元素,通过resultSelector转换成对应的 Observable 或 Promise,然后把这一系列序列**按顺序串联(concatenate)**成一个最终的 Observable。本文以 doc/api/core/operators/for.md 为核心,结合仓库中的 源码实现 与 单元测试,完整讲解该操作符的参数语义、底层串联原理、Promise 兼容、异常传播,以及它和catch、concat等相关操作符的关系,帮助你在实际项目中安全、高效地使用这一"数组驱动"式序列编排能力。
一、API 总览与签名
Rx.Observable.for是一个静态方法,官方文档给出如下签名:
Rx.Observable.for(sources, resultSelector, [thisArg])其语义为:通过依次对sources中的每个元素调用resultSelector得到若干 Observable 序列(或 Promise),并将这些序列按顺序连接成一个输出序列。当某个序列正常结束时,才开始订阅下一个序列,直到所有序列都完成,最终触发onCompleted。
该方法的别名为forIn,二者完全等价,其存在是为了兼容 IE9 以下浏览器(因为for是保留字,在老式浏览器中无法安全地通过Observable.for(...)形式访问,可改用Observable.forIn(...))。
参数详解
| 参数 | 类型 | 必填 | 说明 |
|---|---|---|---|
sources | Array | 是 | 待处理的元素数组,将被逐一转换成可观察序列 |
resultSelector | Function | 是 | 将数组元素映射为 Observable 或 Promise 的函数,调用时依次传入(value, index, sourceArray)三个实参 |
thisArg | Any | 否 | 执行resultSelector时this的指向对象 |
其中resultSelector被调用时携带三个参数:
- value:当前元素的值;
- index:当前元素在数组中的下标;
- sources:正在被遍历的整个源数组(即被订阅的集合对象)。
返回值
(Observable):一个由所有子序列按顺序拼接而成的 Observable。每一个子序列既可以是 Observable,也可以是 Promise(内部会自动将 Promise 包装为 Observable)。
二、官方示例:从数组到串联序列
2.1 使用 Observable 作为映射结果
/* Using Observables */ var array = [1, 2, 3]; var source = Rx.Observable.for( array, function (x) { return Rx.Observable.return(x); }); var subscription = source.subscribe( function (x) { console.log('Next: ' + x); }, function (err) { console.log('Error: ' + err); }, function () { console.log('Completed'); }); // => Next: 1 // => Next: 2 // => Next: 3 // => Completed这里Rx.Observable.return(x)会创建一个只发出一个值x就立即完成的单元素序列,因此for把三个单元素序列串起来,最终依次输出1、2、3并完成。
2.2 使用 Promise 作为映射结果
/* Using Promises */ var array = [1, 2, 3]; var source = Rx.Observable.for( array, function (x) { return RSVP.Promise.resolve(x); }); var subscription = source.subscribe( function (x) { console.log('Next: ' + x); }, function (err) { console.log('Error: ' + err); }, function () { console.log('Completed'); }); // => Next: 1 // => Next: 2 // => Next: 3 // => Completed当resultSelector返回 Promise 时,for同样能够处理——内部会在订阅前检测到 Promise 并自动将其转换为 Observable(详见下一节源码剖析),因此这里以RSVP.Promise.resolve(x)为例同样按顺序输出1、2、3并完成。
三、源码剖析:for的底层实现
3.1 入口实现:一行代码串联
仓库中该操作符的完整实现位于 src/core/linq/observable/for.js:
Observable['for'] = Observable.forIn = function (sources, resultSelector, thisArg) { return enumerableOf(sources, resultSelector, thisArg).concat(); };实现非常精简:先通过enumerableOf(即Rx.internals.Enumerable.of)把(sources, resultSelector, thisArg)包装成一个惰性枚举器(lazy enumerable),再调用其concat()方法得到最终的串联 Observable。注意这里使用方括号写法Observable['for']并同时赋值给forIn,正是为了避免for作为关键字带来的语法问题,这也印证了文档中"IE9 以下使用forIn别名"的说明。
3.2 惰性枚举器:Enumerable.of与三参数回调
enumerableOf定义于 src/core/enumerable.js:
var OfEnumerable = (function(__super__) { inherits(OfEnumerable, __super__); function OfEnumerable(s, fn, thisArg) { this.s = s; this.fn = fn ? bindCallback(fn, thisArg, 3) : null; } OfEnumerable.prototype[$iterator$] = function () { return new OfEnumerator(this); }; function OfEnumerator(p) { this.i = -1; this.s = p.s; this.l = this.s.length; this.fn = p.fn; } OfEnumerator.prototype.next = function () { return ++this.i < this.l ? { done: false, value: !this.fn ? this.s[this.i] : this.fn(this.s[this.i], this.i, this.s) } : doneEnumerator; }; return OfEnumerable; }(Enumerable)); var enumerableOf = Enumerable.of = function (source, selector, thisArg) { return new OfEnumerable(source, selector, thisArg); };关键点:
- 惰性求值:
OfEnumerable只是保存了源数组和回调,只有调用next()时才会真正执行resultSelector; - 三参数调用:
this.fn(this.s[this.i], this.i, this.s)与文档描述的(value, index, sources)完全一致; thisArg绑定:通过bindCallback(fn, thisArg, 3)实现,其定义见 src/core/internal/bindcallback.js。当thisArg未传(undefined)时直接返回原函数;传入时则生成一个把this绑定到thisArg的包装函数,保证resultSelector内部this指向正确;- 可选回调:当
fn为空时(!this.fn),枚举器直接产出源数组元素本身,相当于"只遍历不映射"。
3.3 串联调度:ConcatEnumerableObservable
Enumerable.prototype.concat实现在 src/core/enumerable.js,它返回ConcatEnumerableObservable,其subscribeCore的核心逻辑是:
- 通过
SerialDisposable管理当前正在订阅的子序列,保证同一时刻只订阅一个; - 使用
currentThreadScheduler.scheduleRecursive进行递归调度,配合内部InnerObserver:子序列每完成一次,就通过recurse继续拉取下一个元素并订阅; - 枚举结束(
done === true)时向观察者发出onCompleted; - Promise 自动转换:
isPromise(currentValue) && (currentValue = observableFromPromise(currentValue)),即每个元素在订阅前若检测到是 Promise(借助 src/core/headers/experimentalheader.js 中引入的isPromise与observableFromPromise = Observable.fromPromise),会先被转换为 Observable 再订阅——这正是第二节中 Promise 示例能够直接工作的原因; - 返回值是一个
NAryDisposable(组合了序列订阅、调度句柄与一个IsDisposedDisposable状态标记),因此取消订阅可以一次性释放所有内部资源。
由此可以总结for的整体执行流程:
sources 数组 │ 惰性遍历(每次 next() 取一个元素) ▼ resultSelector(value, index, sources) ──► Observable / Promise │ Promise 先自动包装成 Observable ▼ 逐一订阅(SerialDisposable 保证串行、前一序列完成后才订阅下一个) │ ▼ 所有序列完成 ──► onCompleted;任一序列出错 ──► onError 并停止四、测试验证:行为与异常语义
仓库提供了完整的单元测试 tests/observable/for.js,基于TestScheduler与ReactiveTest断言时序,可直接反应该操作符的实际行为。
4.1 串联顺序与"等待完成"语义(for basic)
测试中为数组[1, 2, 3]的每个元素创建一个"冷序列":
return Observable'for' { return scheduler.createColdObservable( onNext(x * 100 + 10, x * 10 + 1), onNext(x * 100 + 20, x * 10 + 2), onNext(x * 100 + 30, x * 10 + 3), onCompleted(x * 100 + 40) ); });断言结果清晰地展示了串联的时序特征:
onNext(310, 11), onNext(320, 12), onNext(330, 13) // 第 1 个序列(x=1) onNext(550, 21), onNext(560, 22), onNext(570, 23) // 第 2 个序列(x=2),从 550 才开始 onNext(890, 31), onNext(900, 32), onNext(910, 33) // 第 3 个序列(x=3),从 890 才开始 onCompleted(920)可以观察到:第 2、3 个序列并非与第 1 个同时并发,而是分别等到前一个序列在时间点440、780完成之后才启动——这正是**顺序串联(concat)**而非并行合并(merge)的铁证。
4.2 异常传播(for throws)
var error = new Error(); var results = scheduler.startScheduler(function () { return Observable'for' { throw error; }); }); results.messages.assertEqual(onError(200, error));当resultSelector在求值时抛出异常,for不会吞掉错误,而是立即以该异常向订阅者发出onError并终止。从实现上看,这与 src/core/enumerable.js 中tryCatch(state.e.next).call(state.e)的防御式调用一致:一旦枚举器next()抛错,就直接走onError分支。
五、进阶使用与相关操作符对比
5.1 使用thisArg控制回调上下文
当resultSelector依赖某个对象作为this时,传入第三个参数即可:
var ctx = { prefix: 'item-' }; var source = Rx.Observable.for( ['a', 'b'], function (x, i) { // 此处的 this 即 ctx return Rx.Observable.return(this.prefix + i + '-' + x); }, ctx);由 src/core/internal/bindcallback.js 可知,bindCallback(fn, thisArg, 3)会生成function(value, index, collection) { return func.call(thisArg, value, index, collection); }这样的包装,因此thisArg只影响resultSelector的this,而参数顺序始终是(value, index, sources)。
5.2 与Rx.Observable.catch(catchError)的关系
for与catch共享同一套"枚举 + 递归调度"基础设施:Observable.catch(catchError)的实现同样是enumerableOf(items).catchError()(见 src/core/linq/observable/catch.js),对应的CatchErrorObservable就在ConcatEnumerableObservable的隔壁(src/core/enumerable.js)。
二者最大的区别在于错误处理策略:
for(concat语义):任一子序列出错即整体onError终止;catch(catchError语义):某个序列出错后继续尝试下一个序列,最后只报告"最后一个"错误。
因此,如果你需要"数组元素逐个尝试、失败则跳过继续"的容错语义,应当使用Observable.catch而非for。
5.3 与concat/concatAll的对比
for本质上是"数组 + 映射函数 + 惰性求值"与"串行连接"的组合:
- 若你已经拥有一个 Observable 数组,直接用
Observable.concat(array)(concatAll)即可完成串联; for的价值在于按需映射:resultSelector在每个元素被真正订阅前才执行(惰性),且能拿到(value, index, sources)三个上下文参数,适合"根据下标生成不同序列"这类场景;for还额外提供了 Promise 自动包装能力,让数组元素可以直接映射为 Promise 而无需手动fromPromise。
5.4 典型应用场景
- 顺序执行一批异步任务:例如按序请求 N 个接口(每个元素对应一个请求 Promise),且要求上一个请求完成后再发下一个,天然符合
for的串联语义; - 基于下标生成差异化序列:借助
index参数构造延迟时间、重试次数等各不相同的序列; - 作为聚合/实验模块中的基础积木:
for在源码中被归入实验性(experimental)能力,与catch、defer、AsyncSubject等一同封装在 src/core/headers/experimentalheader.js 依赖集合中,适合在完整版rx.all或实验版中直接使用。
六、获取方式与使用前提
for位于实验性(Experimental)能力集合中,从仓库源码结构与文档说明看,可通过以下途径获得:
| 分发形式 | 说明 |
|---|---|
rx.all.js/rx.all.compat.js | 完整版,已包含for |
rx.experimental.js | 实验版,使用时需先加载基础库rx.js/rx.compat.js/rx.lite.js/rx.lite.compat.js之一 |
NPM 包rx | 官方 NPM 发行包 |
NuGetRxJS-Complete/RxJS-Experimental | .NET 生态下的分发包 |
源码位置:
- 实现:src/core/linq/observable/for.js
- 枚举与串联基础设施:src/core/enumerable.js
thisArg绑定:src/core/internal/bindcallback.js- 单元测试:tests/observable/for.js
使用前提:本仓库对应RxJS v4(文档开头即注明 "This is RxJS v 4",最新版本请参考 ReactiveX 官方 RxJS 项目);for依赖实验性模块的头部依赖(isPromise、observableFromPromise、currentThreadScheduler、SerialDisposable等),因此在浏览器环境中请确认加载了正确的构建产物;在 IE9 以下浏览器中建议改用forIn别名。
小结
Rx.Observable.for用极简的 API 完成了"数组 → 序列集合 → 串行拼接"的完整链路:参数层面支持(value, index, sources)三参数回调与可选的thisArg;实现层面由Enumerable.of惰性求值 +ConcatEnumerableObservable递归串行订阅构成,并自动兼容 Promise;测试层面则用TestScheduler精确验证了顺序时序与异常传播。掌握它,你就掌握了 RxJS 中"批量编排顺序任务"的一种基础而可靠的手段。
【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考