尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

RxJS v4 `flatMapWithMaxConcurrent` 操作符完全指南:限流扁平映射与并发控制实战

RxJS v4 `flatMapWithMaxConcurrent` 操作符完全指南:限流扁平映射与并发控制实战 后端【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址https://gitcode.com/gh_mirrors/rxj/RxJS点击查看免费下载本文围绕 RxJS v4The Reactive Extensions for JavaScript中的flatMapWithMaxConcurrent/selectWithMaxConcurrent操作符展开深入讲解其签名、参数语义、返回值并结合本仓库的源码实现与 TypeScript 声明剖析其底层并发控制机制。读完本文你将掌握如何在 RxJS v4 中使用该操作符对可观察序列执行带最大并发上限的扁平映射正确处理 Observable、Promise 与数组/可迭代对象三种输入形态并理解maxConcurrent参数在合并merge环节中如何真正生效。操作符定位与别名关系flatMapWithMaxConcurrent是 RxJS v4 中flatMapselectMany家族的一员它属于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可以对「外层源元素 内层序列元素 各自索引」做进一步组合变换。参数详解参数类型必选说明maxConcurrentNumber是同时被订阅的内部可观察序列的最大数量。当活跃的内部序列数量达到该上限时新产生的内部序列进入等待队列selectorFunction|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限制并发为 2var 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 只有在前一个完成后才会被订阅等待非常适合限制并发、保护下游资源。示例三投影为数组并结合 resultSelectorvar 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)分别对应外层元素值、内层元素值、外层索引、内层索引。例如第一个源元素x1索引ix0投影为[1, 0]两个内层元素依次与源元素做x y ix iy运算得到11002与10012。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 支持投影为ObservableOrPromiseTResult或ArrayOrIterableTResult并分别提供有无resultSelector的重载flatMapWithMaxConcurrentTResult(maxConcurrent: number, selector: _ValueOrSelectorT, ObservableOrPromiseTResult): ObservableTResult; flatMapWithMaxConcurrentTResult(maxConcurrent: number, selector: _ValueOrSelectorT, ArrayOrIterableTResult): ObservableTResult; flatMapWithMaxConcurrentTOther, TResult(maxConcurrent: number, selector: _ValueOrSelectorT, ObservableOrPromiseTOther, resultSelector: special._FlatMapResultSelectorT, TOther, TResult, thisArg?: any): ObservableTResult; flatMapWithMaxConcurrentTOther, TResult(maxConcurrent: number, selector: _ValueOrSelectorT, ArrayOrIterableTOther, resultSelector: special._FlatMapResultSelectorT, TOther, TResult, thisArg?: any): ObservableTResult;同样地selectManyWithMaxConcurrent在 同文件 中提供了完全对称的重载方便使用selectMany命名的项目迁移。文件末尾还附带类型级使用示例覆盖了 Observable、Promise、数组三种 selector 形态以及有无 resultSelector 的六种组合。获取与使用途径在发布产物层面该操作符随以下构建文件分发完整版dist/rx.all.js、dist/rx.all.compat.jscompat 版为兼容旧浏览器的转换版本实验特性版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可做外层与内层元素的联合变换回调参数顺序固定为「外层值、内层值、外层索引、内层索引」使用前务必对照确认。赞分享后端【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址https://gitcode.com/gh_mirrors/rxj/RxJS点击查看免费下载相关推荐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),仅供参考
返回列表