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

资讯详情

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

RxJS(v4)expand 操作符完全指南:递归展开 Observable 的源码剖析与实战

RxJS(v4)expand 操作符完全指南:递归展开 Observable 的源码剖析与实战 后端【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址https://gitcode.com/gh_mirrors/rxj/RxJS点击查看免费下载expand是 RxJS 4Rx.Observable.prototype.expand(selector, [scheduler])中一个极具表现力的递归操作符它把每个产生的元素再次喂给 selector得到新的序列后继续递归展开从而构建出由当前状态推导下一状态的无限/有限推演流程非常适合模拟状态机、广度优先搜索、数论序列生成等场景。读完本文你将掌握 expand 的完整 API 语义、其基于调度器蹦床trampoline的底层实现原理、如何结合take等操作符控制递归终止以及如何用仓库自带的虚拟时间测试验证展开时序。一、expand 是什么递归地展开序列expand的官方定义是Expands an observable sequence by recursively invoking selector通过递归调用 selector 来展开一个可观察序列。它与selectManyflatMap的最大区别在于selectMany只做一层映射拍平而expand会把映射产生的序列中的每个元素再次作为输入去调用 selector如此反复直到序列枯竭或外部截断。从仓库根目录的实现文件 src/core/linq/observable/expand.js 可以看到它的定义与注册方式/** * Expands an observable sequence by recursively invoking selector. * * param {Function} selector Selector function to invoke for each produced element, resulting in another sequence to which the selector will be invoked recursively again. * param {Scheduler} [scheduler] Scheduler on which to perform the expansion. If not provided, this defaults to the current thread scheduler. * returns {Observable} An observable sequence containing all the elements produced by the recursive expansion. */ observableProto.expand function (selector, scheduler) { isScheduler(scheduler) || (scheduler currentThreadScheduler); return new ExpandObservable(this, selector, scheduler); };注意两个与 API 文档的细微出入这里以源码为准默认调度器API 文档写的是默认 immediate 调度器但源码与源码注释一致表明默认值是currentThreadScheduler当前线程调度器。它通过isScheduler(scheduler) || (scheduler currentThreadScheduler)完成兜底。返回值API 文档中的返回值描述判定是否全部元素通过谓词明显是从其他操作符如all/every复制粘贴而来与实现不符。依据 src/core/linq/observable/expand.js 的注释正确语义是返回包含递归展开所产生的全部元素的 Observable 序列。二、API 签名与参数说明Rx.Observable.prototype.expand(selector, [scheduler])参数类型必填说明selectorFunction是对每个产生的元素调用的选择器函数返回一个新的序列该序列中的每个元素又会被递归地交给 selector 继续展开。签名形如function (x) { return observable; }schedulerScheduler否执行展开操作的调度器。不传时默认使用当前线程调度器Rx.Scheduler.currentThread即Scheduler.currentThread new CurrentThreadScheduler()见 src/core/concurrency/currentthreadscheduler.js返回Observable—— 包含递归展开产生的所有元素的序列。三、官方示例从 42 开始的倍增展开API 文档给出了一个最直观的例子以42为种子值每次把当前值加 42 得到新序列递归展开后用take(5)截断前 5 个元素var source Rx.Observable.return(42) .expand(function (x) { return Rx.Observable.return(42 x); }) .take(5); var subscription source.subscribe( function (x) { console.log(Next: %s, x); }, function (err) { console.log(Error: %s, err); }, function () { console.log(Completed); }); // Next: 42 // Next: 84 // Next: 126 // Next: 168 // Next: 210 // Completed运行结果输出42, 84, 126, 168, 210后正常Completed。这个例子揭示了两个关键点种子来自源序列第一个元素42由源序列Rx.Observable.return(42)产生之后每个元素x都会触发一次selector(x)无限推演必须外部截断expand本身是无限递归的每次42 x都会产生新元素因此必须搭配take(5)之类的前置终止手段否则订阅将永不完成。这一点可以在测试expand never中印证——永不结束的源会让订阅一直保持见下文第五节。四、源码级原理队列驱动 调度器蹦床的递归展开expand的实现由两个类组成src/core/linq/observable/expand.jsExpandObservable继承ObservableBase负责管理待展开序列的队列与调度ExpandObserver继承AbstractObserver负责消费每个元素并产出下一轮序列。4.1 状态容器队列、计数与资源管理subscribeCore中构造了递归展开的共享状态src/core/linq/observable/expand.jsExpandObservable.prototype.subscribeCore function (o) { var m new SerialDisposable(), d new CompositeDisposable(m), state { q: [], // 待展开的 observable 队列广度优先 m: m, // 串行占位 disposable指向当前递归调度 d: d, // 复合 disposable聚合所有子订阅 activeCount: 0, // 仍在进行中的子序列数量 isAcquired: false, // 是否已有调度器蹦床在工作 o: o // 下游观察者 }; state.q.push(this.source); state.activeCount; this._ensureActive(state); return d; };设计要点队列q保证广度优先展开顺序源序列入队后每个元素展开出的新序列都被追加到队尾scheduleRecursive每次从队首取出一个序列订阅形成层一层扩散的效果避免深度递归导致调用栈溢出activeCount控制终止只有所有进行中的子序列都完成activeCount 0时才向下游发出onCompleted资源管理CompositeDisposablesrc/core/disposables/compositedisposable.js聚合所有子订阅整体退订时一并释放SerialDisposablesrc/core/disposables/booleandisposable.js占位管理递归调度本身。4.2 调度器蹦床避免递归栈溢出真正的递归不是用 JS 函数递归而是交给调度器的scheduleRecursive蹦床机制ExpandObservable.prototype._ensureActive function (state) { var isOwner false; if (state.q.length 0) { isOwner !state.isAcquired; state.isAcquired true; } isOwner state.m.setDisposable(this._scheduler.scheduleRecursive([state, this], scheduleRecursive)); }; function scheduleRecursive(args, recurse) { var state args[0], self args[1]; var work; if (state.q.length 0) { work state.q.shift(); } else { state.isAcquired false; return; } var m1 new SingleAssignmentDisposable(); state.d.add(m1); m1.setDisposable(work.subscribe(new ExpandObserver(state, self, m1))); recurse([state, self]); }isAcquired是一个所有权标记同一时刻只允许一个蹦床循环在工作新元素到来时若已有循环在跑就只把新序列追加进队列由现有循环继续消费队列耗尽时重置isAcquired false等待下一轮元素到来时再触发新的蹦床scheduleRecursive由 src/core/concurrency/scheduler.recursive.js 提供它把每次递归调用重新放回调度器队列执行默认的CurrentThreadScheduler使用优先级队列蹦床见 src/core/concurrency/currentthreadscheduler.js从而把递归摊平成迭代从根本上规避深层递归的调用栈溢出风险。4.3 ExpandObserver元素的三条出路每个被订阅的子序列都用ExpandObserver包装src/core/linq/observable/expand.jsExpandObserver.prototype.next function (x) { this._s.o.onNext(x); // 1. 元素直接转发给下游 var result tryCatch(this._p._fn)(x); // 2. 调用 selector 产出新序列 if (result errorObj) { return this._s.o.onError(result.e); } this._s.q.push(result); // 3. 新序列入队继续展开 this._s.activeCount; this._p._ensureActive(this._s); }; ExpandObserver.prototype.error function (e) { this._s.o.onError(e); // 错误原样转发终止 }; ExpandObserver.prototype.completed function () { this._s.d.remove(this._m1); this._s.activeCount--; this._s.activeCount 0 this._s.o.onCompleted(); // 全部结束才完成 };next先向订阅者推送元素再用tryCatch包裹 selector 调用selector 抛错立即转成onErrortryCatch/errorObj来自仓库 internal 基础设施参见 src/core/internal 目录error任何子序列出错错误直接透传并终止整个展开completed单个子序列完成只减少计数只有activeCount 0时才通知下游onCompleted——这是递归展开整体完成的正确判定。4.4 与 ObservableBase 的衔接ExpandObservable继承自ObservableBasesrc/core/perf/observablebase.js后者负责统一的_subscribe封装在默认当前线程调度器上调度subscribeCore并通过AutoDetachObserversrc/core/autodetachobserver.js实现下游观察者抛错时的自动退订与异常重抛从而保证expand与其余操作符在订阅/退订语义上完全一致。五、用测试用例验证展开语义仓库在 tests/observable/expand.js 中用 QUnit TestScheduler虚拟时间完整覆盖了展开的五个分支行为是理解语义的最佳佐证测试用例场景期望行为expand empty源序列为空300ms 完成直接onCompleted(300)且 selector 从未被调用expand error源序列出错错误onError(300, error)原样透传订阅随即结束expand never源序列永不结束无任何消息输出订阅从 201ms 保持到 1000ms测试窗口结束expand basic常规递归展开展开产生的子序列元素按广度优先顺序依次下发expand throwselector 抛异常元素仍先推送随后立刻onError(550, error)其中expand basic最有价值源在 550ms、850ms 各推1、2selector 对x返回一个100ms 推2x、200ms 推3x、300ms 完成的冷序列。最终结果序列1, 2, 3, 4, 2, 6, 6, 8, 4, 9, 12, 12, 12, 16完整呈现了广度优先的展开轨迹1 → 2、32、3 继续展开 → 4、6、6、8……。expand throw则验证了先推元素、再报 selector 错误的精确时序。六、使用注意与最佳实践必须考虑终止条件expand天然是无限推演。实际使用中几乎总是配合take(n)截断数量、takeWhile按谓词截断或让 selector 在满足条件时返回空序列如Rx.Observable.empty()来结束递归否则订阅永不完成对应expand never的行为。默认调度器是当前线程调度器默认情况下展开过程在订阅线程同步蹦床执行需要异步化时显式传入Rx.Scheduler.default、Rx.Scheduler.timeout等调度器。错误处理selector 内抛错或任一子序列报错都会立即终止整个展开并把错误传给下游onError无需额外防御。内存与资源展开的每一层都会产生订阅整体由CompositeDisposable管理及时退订dispose可一次性释放所有层级的资源。典型场景状态机推演由当前状态生成后继状态序列、广度优先搜索/树的逐层遍历、迭代公式序列生成等由旧值产生新值、新值继续参与计算的递归模型。七、获取与使用 expandexpand位于仓库的 experimental 功能集其实现挂在observableProto上随构建产物分发源码位置src/core/linq/observable/expand.js由src/core/linq/observable目录下各操作符模块共同构成核心库模块产物仓库 modules 目录下提供按功能拆分的 npm 包其中rx-lite-experimental及rx-lite-experimental-compat即包含 expand 等实验性操作符的构建产物NuGet对应RxJS-All完整包与RxJS-Experimental实验功能包NPM发布为rx包v4 系列需配合基础运行时rx.js/rx.compat.js/rx.lite.js/rx.lite.compat.js等前置模块使用详见 doc/api/core/operators/expand.md 的 Location 章节测试参考tests/observable/expand.js 可作为接入与验证的样板使用Rx.TestSchedulerRx.ReactiveTest断言消息时序。总结expand通过元素转发 selector 递归 调度器蹦床 队列广度优先的组合把递归展开从危险的深度调用栈问题中解放出来使其成为安全、可预测的推演型操作符。理解它的队列状态机、activeCount终止判定与默认当前线程调度器这三个核心机制你就能在状态推演与逐层遍历类问题中放心地使用它并借助take系列操作符精确控制展开的边界。赞分享后端【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址https://gitcode.com/gh_mirrors/rxj/RxJS点击查看免费下载相关推荐RxJS v4 combineLatest 操作符完全指南API 用法、行为语义与源码实现剖析RxJS v4 combineLatest 操作符完全指南API 用法、行为语义与源码实现剖析 本文基于 RxJS v4The Reactive Exten后端RxJava 自定义 Observable 操作符完全指南lift 序列操作符与 compose 转换操作符的源码级实战RxJava 自定义 Observable 操作符完全指南lift 序列操作符与 compose 转换操作符的源码级实战 在 RxJava 中编写自定义 Ob后端异步编程yq flatten 操作符完全指南递归展平嵌套数组的原理与实战yq flatten 操作符完全指南递归展平嵌套数组的原理与实战 flatten 是 yq 中用于将嵌套数组递归展平的专用操作符它能把多层嵌套的序列结构拍平开发工具CLI创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表