- 后端
【免费下载链接】RxJS
The Reactive Extensions for JavaScript
本文围绕 RxJS v4(The Reactive Extensions for JavaScript)中的flatMapWithMaxConcurrent/selectWithMaxConcurrent操作符展开,深入讲解其签名、参数语义、返回值,并结合本仓库的源码实现与 TypeScript 声明剖析其底层并发控制机制。读完本文,你将掌握如何在 RxJS v4 中使用该操作符对可观察序列执行带最大并发上限的扁平映射,正确处理 Observable、Promise 与数组/可迭代对象三种输入形态,并理解maxConcurrent参数在合并(merge)环节中如何真正生效。
操作符定位与别名关系
flatMapWithMaxConcurrent是 RxJS v4 中flatMap(selectMany)家族的一员,它属于Rx.Observable.prototype上的实例方法。官方文档给出的完整签名为:
Rx.Observable.prototype.flatMapWithMaxConcurrent(maxConcurrent, selector, [resultSelector], [thisArg]) Rx.Observable.prototype.selectWithMaxConcurrent(maxConcurrent, selector, [resultSelector], [thisArg])其中flatMapWithMaxConcurrent与selectWithMaxConcurrent互为别名,二者指向同一实现。若使用旧式select命名风格,还可写作selectManyWithMaxConcurrent(见 TypeScript 声明文件)。
需要特别说明的是,仓库源码中还注册了第三个等价名称flatMapMaxConcurrent。在 src/core/perf/operators/flatmapwithmaxconcurrent.js 中可以看到:
observableProto.flatMapWithMaxConcurrent = observableProto.flatMapMaxConcurrent = function(limit, selector, resultSelector, thisArg) { return new FlatMapObservable(this, selector, resultSelector, thisArg).merge(limit); };这段实现非常简洁,却揭示了该操作符的本质:先通过FlatMapObservable完成「一对多」投影(flatMap 阶段),再调用带limit参数的merge完成「限流合并」(merge 阶段)。因此,理解flatMapWithMaxConcurrent的关键,就是分别理解这两个底层组件——我们会在后文逐一拆解。
功能语义:并发受限的扁平映射
该操作符将源可观察序列中的每个元素投影(project)为另一个可观察序列,然后将产生的这些内部序列合并为一个输出序列;与普通flatMap不同的是,它最多同时订阅maxConcurrent个内部序列,超出上限的内部序列会被排队等待,直到有内部序列完成后再依次订阅。
具体来说,它支持三种投影形态:
- 投影为 Observable:将每个源元素映射为一个可观察序列,再合并其结果;
- 投影为 Promise:将每个源元素映射为 Promise,内部自动通过
fromPromise转换为可观察序列; - 投影为数组/可迭代对象:将每个源元素映射为数组或可迭代对象,内部自动展开为可观察序列。
同时,可选的结果选择器resultSelector可以对「外层源元素 + 内层序列元素 + 各自索引」做进一步组合变换。
参数详解
| 参数 | 类型 | 必选 | 说明 |
|---|---|---|---|
maxConcurrent | Number | 是 | 同时被订阅的内部可观察序列的最大数量。当活跃的内部序列数量达到该上限时,新产生的内部序列进入等待队列 |
selector | Function|Iterable|Promise | 是 | 投影目标或变换函数。可以是一个把每个元素变换为序列的函数,也可以直接是一个 Observable / Promise / 数组 / 可迭代对象(此时所有源元素都被投影到同一个序列上)。作为函数调用时依次接收:(1)元素的值、(2)元素的索引、(3)正在被订阅的 Observable 对象 |
[resultSelector] | Function | 否 | 对中间序列每个元素施加的变换函数,依次接收:(1)外层元素的值、(2)内层元素的值、(3)外层元素的索引、(4)内层元素的索引 |
[thisArg] | Any | 否 | 当resultSelector不是函数时,作为执行selector时this的绑定对象 |
关于thisArg的语义需要留意:文档明确指出它仅在resultSelector不是函数时生效,作为selector执行时的this上下文。
返回值
返回一个Observable:其元素是对输入序列每个元素调用「一对多」投影函数后,再把每个投影出的序列元素与其对应的源元素共同映射(resultSelector)得到的最终结果。也就是说,返回值仍然是可观察序列,可直接subscribe或继续链式调用其他操作符。
完整示例
示例一:投影为 Observable(限制并发为 2)
var source = Rx.Observable.range(0, 5) .flatMapWithMaxConcurrent(2, function (x, i) { return Rx.Observable .interval(100) .take(x).map(function() { return i; }); }); var subscription = source.subscribe( function (x) { console.log('Next: %s', x); }, function (err) { console.log('Error: %s', err); }, function () { console.log('Completed'); }); // => Next: 1 // => Next: 2 // => Next: 3 // => Next: 2 // => Next: 3 // => Next: 4 // => Next: 3 // => Next: 4 // => Next: 4 // => Next: 4 // => Completed这个示例最能体现并发限制的直观效果:maxConcurrent = 2意味着同一时刻最多只有两个interval内部序列处于活跃状态,其余内部序列等待前序序列结束后才被订阅,因此输出呈现「两两交错推进」而非全部并发的形态。
示例二:投影为 Promise(并发为 1,即串行)
var source = Rx.Observable.of(1,2,3,4) .flatMapWithMaxConcurrent(1, function (x, i) { return Promise.resolve(x + i); }); var subscription = source.subscribe( function (x) { console.log('Next: %s', x); }, function (err) { console.log('Error: %s', err); }, function () { console.log('Completed'); }); // => Next: 1 // => Next: 3 // => Next: 5 // => Next: 7 // => Completed当maxConcurrent = 1时,该操作符退化为严格的串行执行(等价于 concatMap 的语义):每个 Promise 只有在前一个完成后才会被订阅等待,非常适合限制并发、保护下游资源。
示例三:投影为数组并结合 resultSelector
var source = Rx.Observable.of(1,2,3) .flatMapWithMaxConcurrent( 1, function (x, i) { return [x,i]; }, function (x, y, ix, iy) { return x + y + ix + iy; } ); var subscription = source.subscribe( function (x) { console.log('Next: %s', x); }, function (err) { console.log('Error: %s', err); }, function () { console.log('Completed'); }); // => Next: 2 // => Next: 2 // => Next: 5 // => Next: 5 // => Next: 8 // => Next: 8 // => Completed该示例演示了resultSelector的四参数形态:(x, y, ix, iy)分别对应外层元素值、内层元素值、外层索引、内层索引。例如第一个源元素x=1(索引ix=0)投影为[1, 0],两个内层元素依次与源元素做x + y + ix + iy运算,得到1+1+0+0=2与1+0+0+1=2。
selector 直接传序列的第三种用法
文档还给出了 selector 直接传 Observable / Promise / 数组的形态(此时所有源元素都被投影到同一个序列上):
source.flatMapWithMaxConcurrent(1, Rx.Observable.of(1,2,3)); source.flatMapWithMaxConcurrent(1, Promise.resolve(42)); source.flatMapWithMaxConcurrent(1, [1,2,3]);源码级原理解析
第一阶段:FlatMapObservable 完成投影归一化
FlatMapObservable定义于 src/core/perf/operators/flatmapbase.js,其构造函数与核心逻辑如下:
function FlatMapObservable(source, selector, resultSelector, thisArg) { this.resultSelector = isFunction(resultSelector) ? resultSelector : null; this.selector = bindCallback(isFunction(selector) ? selector : function() { return selector; }, thisArg, 3); this.source = source; __super__.call(this); }关键点有两个:
- 若
selector不是函数(即直接传入了 Observable / Promise / 数组),源码会用function() { return selector; }将其包装成恒定返回该对象的函数,从而统一处理三种输入形态; - 对
selector执行bindCallback(..., thisArg, 3),将thisArg绑定为this,并约定回调最多接收 3 个参数(值、索引、源 Observable)。
真正逐元素投影的逻辑在InnerObserver.next中:
InnerObserver.prototype.next = function(x) { var i = this.i++; var result = tryCatch(this.selector)(x, i, this.source); if (result === errorObj) { return this.o.onError(result.e); } isPromise(result) && (result = observableFromPromise(result)); (isArrayLike(result) || isIterable(result)) && (result = Observable.from(result)); this.o.onNext(this._wrapResult(result, x, i)); };这里可以看到完整的归一化链路:selector 抛错 → 立即 onError;返回 Promise → 通过observableFromPromise转为 Observable;返回数组/类数组/可迭代对象 → 通过Observable.from展开为 Observable。此外,若提供了resultSelector,_wrapResult会对内部序列的每个元素执行resultSelector(x, y, i, i2)(外层值、内层值、外层索引、内层索引)完成二次映射。
第二阶段:MergeObservable 实现限流排队
投影产出的「序列的序列」随后交给带maxConcurrent参数的merge处理,其实现位于 src/core/perf/operators/mergeconcat.js。核心的MergeObserver用三个状态字段完成限流:
function MergeObserver(o, max, g) { this.o = o; this.max = max; this.g = g; // CompositeDisposable,统一管理所有内部订阅 this.done = false; // 源序列是否已完成 this.q = []; // 等待队列:超出上限的内部序列先进先出排队 this.activeCount = 0; // 当前活跃的内部序列数量 }- 每当源序列产出新的内部序列(
next),若activeCount < max则activeCount++并立即订阅;否则推入q队列等待; - 每个内部序列用一个
SingleAssignmentDisposable管理订阅,并注册进CompositeDisposable; - 当某个内部序列完成(
InnerObserver.completed)时,先从g中移除该订阅;若q中还有等待者,则取出队首序列(q.shift())继续订阅——注意此时活跃数不减少,新序列无缝顶替;若队列已空,则activeCount--; - 只有当源序列完成(
done = true)且activeCount === 0时,整个输出序列才onCompleted; - 任一层级出错都会立即
onError终止整个序列。
因此maxConcurrent真正控制的不是「发射速度」而是「内部序列的订阅并发度」,这正是该操作符用于限流(backpressure)场景的根基。
模块化版本
仓库还提供了独立模块化的实现,便于按需引入:入口为 src/modular/observable/flatmapmaxconcurrent.js,其内部同样组合了 src/modular/observable/flatmapobservable.js 与 src/modular/observable/mergeconcat.js,逻辑与src/core/perf下的实现一一对应:
module.exports = function flatMapLatest (source, limit, selector, resultSelector, thisArg) { return mergeConcat(new FlatMapObservable(source, selector, resultSelector, thisArg), limit); };TypeScript 类型声明
在 ts/core/linq/observable/flatmapwithmaxconcurrent.ts 中可以看到完整的重载签名,selector 支持投影为ObservableOrPromise<TResult>或ArrayOrIterable<TResult>,并分别提供有无resultSelector的重载:
flatMapWithMaxConcurrent<TResult>(maxConcurrent: number, selector: _ValueOrSelector<T, ObservableOrPromise<TResult>>): Observable<TResult>; flatMapWithMaxConcurrent<TResult>(maxConcurrent: number, selector: _ValueOrSelector<T, ArrayOrIterable<TResult>>): Observable<TResult>; flatMapWithMaxConcurrent<TOther, TResult>(maxConcurrent: number, selector: _ValueOrSelector<T, ObservableOrPromise<TOther>>, resultSelector: special._FlatMapResultSelector<T, TOther, TResult>, thisArg?: any): Observable<TResult>; flatMapWithMaxConcurrent<TOther, TResult>(maxConcurrent: number, selector: _ValueOrSelector<T, ArrayOrIterable<TOther>>, resultSelector: special._FlatMapResultSelector<T, TOther, TResult>, thisArg?: any): Observable<TResult>;同样地,selectManyWithMaxConcurrent在 同文件 中提供了完全对称的重载,方便使用selectMany命名的项目迁移。文件末尾还附带类型级使用示例,覆盖了 Observable、Promise、数组三种 selector 形态以及有无 resultSelector 的六种组合。
获取与使用途径
在发布产物层面,该操作符随以下构建文件分发:
- 完整版:dist/rx.all.js、dist/rx.all.compat.js(compat 版为兼容旧浏览器的转换版本);
- 实验特性版:dist/rx.experimental.js,同时提供对应的
.map与.min.js变体; - 模块化独立包:modules/rx-lite-experimental/rx.lite.experimental.js(及 compat 版本)。
在 NPM 上,该能力归属于rx包(见 package.json);在 NuGet 上,则包含于RxJS-Complete与RxJS-Experimental两个包中。引入对应构建文件后,即可通过Rx.Observable.prototype.flatMapWithMaxConcurrent直接调用。
实践要点小结
maxConcurrent是订阅并发上限而非发射速率上限:超出限制的内部序列会按到达顺序排队(FIFO),一个内部序列完成后自动从队首取出下一个;maxConcurrent = 1等价于串行:适合需要严格顺序或保护稀缺资源的场景;- 三种输入形态统一处理:selector 返回 Observable、Promise、数组/可迭代对象均可,源码内部完成归一化,无需手动转换;
- 错误传播立即终止:selector 抛错或任一层级出错都会直接
onError,不会继续排队等待; - 配合
resultSelector可做外层与内层元素的联合变换,回调参数顺序固定为「外层值、内层值、外层索引、内层索引」,使用前务必对照确认。
- 后端
【免费下载链接】RxJS
The Reactive Extensions for JavaScript
相关推荐
es-toolkit flatMapAsync 完全指南:异步映射 + 单层扁平化的高效实现与并发控制
es toolkit flatMapAsync 完全指南:异步映射 + 单层扁平化的高效实现与并发控制 导读 flatMapAsync 是 es toolkit
前端后端RxJS v4 `mergeAll` 运算符完全解析:将高阶 Observable 扁平合并为单一序列
RxJS v4 mergeAll 运算符完全解析:将高阶 Observable 扁平合并为单一序列 mergeAll (旧名 mergeObservable )
后端RxJS 高阶 Observable 完全指南:concatMap / mergeMap / switchMap / exhaustMap 扁平化操作符深度解析
RxJS 高阶 Observable 完全指南:concatMap / mergeMap / switchMap / exhaustMap 扁平化操作符深度解析
前端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考