news 2026/9/20 19:39:41

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
RxJS(v4)expand 操作符完全指南:递归展开 Observable 的源码剖析与实战
  • 后端

【免费下载链接】RxJS

The Reactive Extensions for JavaScript

项目地址:https://gitcode.com/gh_mirrors/rxj/RxJS
点击查看免费下载

expand是 RxJS 4(Rx.Observable.prototype.expand(selector, [scheduler]))中一个极具表现力的递归操作符:它把每个产生的元素再次喂给 selector,得到新的序列后继续递归展开,从而构建出"由当前状态推导下一状态"的无限/有限推演流程,非常适合模拟状态机、广度优先搜索、数论序列生成等场景。读完本文,你将掌握 expand 的完整 API 语义、其基于调度器蹦床(trampoline)的底层实现原理、如何结合take等操作符控制递归终止,以及如何用仓库自带的虚拟时间测试验证展开时序。

一、expand 是什么:递归地展开序列

expand的官方定义是:Expands an observable sequence by recursively invoking selector(通过递归调用 selector 来展开一个可观察序列)。它与selectMany(flatMap)的最大区别在于: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 文档的细微出入,这里以源码为准:

  1. 默认调度器:API 文档写的是"默认 immediate 调度器",但源码与源码注释一致表明默认值是currentThreadScheduler(当前线程调度器)。它通过isScheduler(scheduler) || (scheduler = currentThreadScheduler)完成兜底。
  2. 返回值: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.js):

  • ExpandObservable:继承ObservableBase,负责管理待展开序列的队列与调度;
  • ExpandObserver:继承AbstractObserver,负责消费每个元素并产出下一轮序列。

4.1 状态容器:队列、计数与资源管理

subscribeCore中构造了递归展开的共享状态(src/core/linq/observable/expand.js):

ExpandObservable.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
  • 资源管理CompositeDisposable(src/core/disposables/compositedisposable.js)聚合所有子订阅,整体退订时一并释放;SerialDisposable(src/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.js):

ExpandObserver.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继承自ObservableBase(src/core/perf/observablebase.js),后者负责统一的_subscribe封装:在默认当前线程调度器上调度subscribeCore,并通过AutoDetachObserver(src/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 各推12,selector 对x返回一个"100ms 推2x、200ms 推3x、300ms 完成"的冷序列。最终结果序列1, 2, 3, 4, 2, 6, 6, 8, 4, 9, 12, 12, 12, 16完整呈现了广度优先的展开轨迹(1 → 2、3;2、3 继续展开 → 4、6、6、8……)。expand throw则验证了"先推元素、再报 selector 错误"的精确时序。

六、使用注意与最佳实践

  1. 必须考虑终止条件expand天然是无限推演。实际使用中几乎总是配合take(n)(截断数量)、takeWhile(按谓词截断)或让 selector 在满足条件时返回空序列(如Rx.Observable.empty())来结束递归,否则订阅永不完成(对应expand never的行为)。
  2. 默认调度器是当前线程调度器:默认情况下展开过程在订阅线程同步蹦床执行;需要异步化时显式传入Rx.Scheduler.defaultRx.Scheduler.timeout等调度器。
  3. 错误处理:selector 内抛错或任一子序列报错都会立即终止整个展开并把错误传给下游onError,无需额外防御。
  4. 内存与资源:展开的每一层都会产生订阅,整体由CompositeDisposable管理;及时退订(dispose)可一次性释放所有层级的资源。
  5. 典型场景:状态机推演(由当前状态生成后继状态序列)、广度优先搜索/树的逐层遍历、迭代公式序列生成等"由旧值产生新值、新值继续参与计算"的递归模型。

七、获取与使用 expand

expand位于仓库的 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.TestScheduler+Rx.ReactiveTest断言消息时序。

总结

expand通过"元素转发 + selector 递归 + 调度器蹦床 + 队列广度优先"的组合,把递归展开从危险的深度调用栈问题中解放出来,使其成为安全、可预测的推演型操作符。理解它的队列状态机、activeCount终止判定与默认当前线程调度器这三个核心机制,你就能在状态推演与逐层遍历类问题中放心地使用它,并借助take系列操作符精确控制展开的边界。

  • 后端

【免费下载链接】RxJS

The Reactive Extensions for JavaScript

项目地址:https://gitcode.com/gh_mirrors/rxj/RxJS
点击查看免费下载

相关推荐

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/20 19:39:32

在线测速数字背后:带宽、延迟与丢包如何影响你的网速

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/20 19:37:37

自动对局脚本技术拆解:从图像识别到模拟点击的攻防博弈

前阵子在游戏交流群里又看到有人发“向僵尸开炮开挂脚本,小橘子自动对局”这种引流消息,配着几张击杀数截图,下面一群人喊着“求分享”。我盯着那行字看了半天,第一反应不是“这脚本怎么下”,而是“这种自动对局到底是…

作者头像 李华
网站建设 2026/9/20 19:36:46

数字化转型中的数据架构与数据治理:从体系设计到落地实践

简介:这份PPT共43页,是一套基于麦肯锡企业架构(EA)方法论的数字化转型与数据治理规划方案,适合企业架构师、数据管理人员、信息化部门负责人以及咨询项目成员参考。内容以某企业未来五年面临的内外部挑战为背景&#x…

作者头像 李华
网站建设 2026/9/20 19:35:21

SSA-FMD故障诊断工具:MATLAB参数自动寻优与降噪提取

简介:面向具备一定MATLAB编程基础与信号处理知识的研究人员、工程师及高校师生,该方案以机械故障诊断为场景,将麻雀搜索算法(SSA)与特征模态分解(FMD)相结合,用于在强噪声、非平稳条…

作者头像 李华
网站建设 2026/9/20 19:34:43

MediaPipe封装为Windows DLL的完整工程实践

1. 项目概述:为什么要把MediaPipe塞进Windows DLL里MediaPipe这个东西,我最早是在做手势识别demo时接触的,当时用Python跑官方示例,模型推理快、API干净,但一到实际产线部署就卡壳了——客户明确要求“必须是C写的exe&…

作者头像 李华