RxJS 4 的Rx.Observable.repeat深入解析:固定值重复发射、调度器与源码实现
【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS
导读
Rx.Observable.repeat是 RxJS 4(Reactive Extensions for JavaScript)中一个非常基础但高频使用的静态工厂方法:它生成一个可观察序列(Observable),把给定的同一个元素按指定次数重复发射出去。本文以 doc/api/core/operators/repeat.md 为核心,完整讲解它的参数语义、默认行为、无限重复与退订机制,并结合仓库中的性能优化实现、可枚举序列实现与全套单元测试,从 API 用法一路深入到scheduleRecursive递归调度的底层原理,帮助你既会用、又看得懂它的内部机制。
一、API 签名与核心语义
repeat的完整签名如下:
Rx.Observable.repeat(value, [repeatCount], [scheduler])它"生成一个可观察序列,该序列使用指定的调度器(scheduler)发送观察者消息,将给定的元素重复指定次数"。换言之,它产出的序列只包含value这一个值,区别只在于:
- 次数有限:发射
repeatCount次后正常完成(onCompleted); - 次数无限:不指定
repeatCount时无限重复发射,永不完成。
一个直观的对比是它与Rx.Observable.return(即just)的关系:return(value)只发射一次;repeat(value, n)相当于把return的序列"自我复制"了 n 次再串联起来。
二、参数详解与默认值
| 参数 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
value | Any | 是 | 无 | 要重复发射的元素,可以是任意 JS 值(数字、字符串、对象等) |
repeatCount | Number | 否 | -1 | 重复发射的次数;不传或传null时等价于无限重复 |
scheduler | Scheduler | 否 | Scheduler.immediate | 运行"生产者循环"的调度器 |
关于默认调度器,文档标注为Scheduler.immediate(同步执行);值得注意的是,本仓库的性能优化版本 src/core/perf/operators/repeat.js 中实际写法是:
Observable.repeat = function (value, repeatCount, scheduler) { isScheduler(scheduler) || (scheduler = currentThreadScheduler); return new RepeatObservable(value, repeatCount, scheduler); };即在不传调度器时回退到currentThreadScheduler。无论immediate还是currentThreadScheduler,它们都属于"同步执行"调度器,对绝大多数调用方来说行为一致;只有在希望把发射循环推迟到某个队列或异步调度器上执行时,才需要显式传入第三个参数。
关于repeatCount的边界语义(可从 tests/observable/repeat.js 的测试中逐一确认):
repeat(42, 0):一个值都不发射,直接onCompleted;repeat(42, 1):发射 1 次42,然后onCompleted;repeat(42, 10):发射 10 次42,然后onCompleted;repeat(42, -1)/repeat(42):无限发射,只有显式退订才会停止。
三、快速上手:完整的可运行示例
原文示例直接可运行,这里保留完整形态并补充注释:
var source = Rx.Observable.repeat(42, 3); var subscription = source.subscribe( function (x) { console.log('Next: ' + x); }, function (err) { console.log('Error: ' + err); }, function () { console.log('Completed'); }); //=> Next: 42 //=> Next: 42 //=> Next: 42 //=> Completed输出结果与参数严格对应:repeatCount = 3时控制台打印三个Next: 42,随后触发Completed。若把第二参数去掉改为Rx.Observable.repeat(42),则Next: 42会无限打印下去,永远不会走到Completed分支——这正是"无限重复"的真实运行效果。
四、无限重复与退订(Dispose)
无限序列在浏览器控制台或 Node 脚本中运行会一直占用主线程输出。正确的做法是通过订阅返回的subscription主动退订:
var source = Rx.Observable.repeat(42); // 不传次数 = 无限 var count = 0; var subscription = source.subscribe(function (x) { console.log('Next: ' + x); if (++count >= 5) { subscription.dispose(); // 手动停止 } }); //=> Next: 42 (共 5 次后停止)这一点在单元测试中有专门覆盖:tests/observable/repeat.js的'repeat value count dispose'与'repeat value'两个用例,都是通过TestScheduler在指定虚拟时刻(如disposed: 207)终止订阅,从而验证"有限次数未跑完时退订能立即中断"以及"无限重复可通过退订中断"。
五、与原型方法observableProto.repeat的区别
除了静态工厂Rx.Observable.repeat(value, ...),RxJS 4 还提供了对已有序列重复的实例方法:
var source = Rx.Observable.return(1).repeat(3);它的实现位于 src/core/linq/observable/repeatproto.js:
observableProto.repeat = function (repeatCount) { return enumerableRepeat(this, repeatCount).concat(); };这里复用了 src/core/enumerable.js 中的RepeatEnumerable/RepeatEnumerator:把"源序列本身"作为可枚举对象,用concat()将同一序列反复串联。注意两者的语义差异:
- 静态
Observable.repeat(value, n):重复的是单个值; - 原型
source.repeat(n):重复的是整个序列——源序列每完成一次,就重新订阅并重放一遍,直到凑满 n 轮(或无限轮)。
对应地,retry操作符在 src/core/linq/observable/retry.js 中就是enumerableRepeat(this, retryCount).catchError(),可见repeat与"重试"是同一套枚举-串联机制的两种用法。
六、源码级剖析:RepeatObservable与递归调度
性能优化版实现在 src/core/perf/operators/repeat.js,核心是两个对象:
1.RepeatObservable(继承ObservableBase)
function RepeatObservable(value, repeatCount, scheduler) { this.value = value; this.repeatCount = repeatCount == null ? -1 : repeatCount; // null → -1 → 无限 this.scheduler = scheduler; __super__.call(this); }构造函数在这里完成了默认值归一化:repeatCount == null一律转为-1,后续循环逻辑以-1作为"无限"哨兵值。
2.RepeatSink的run方法(发射循环本体)
RepeatSink.prototype.run = function () { var observer = this.observer, value = this.parent.value; function loopRecursive(i, recurse) { if (i === -1 || i > 0) { observer.onNext(value); i > 0 && i--; } if (i === 0) { return observer.onCompleted(); } recurse(i); } return this.parent.scheduler.scheduleRecursive(this.parent.repeatCount, loopRecursive); };这段代码是整个操作符的灵魂,逐行解读:
- 循环携带状态
i(剩余发射次数),初值为repeatCount; i === -1(无限)或i > 0(还有剩余)时发射一次value,有限模式下同时自减;i === 0时调用observer.onCompleted()终止整个序列;- 否则通过回调参数
recurse(i)触发下一轮。
scheduleRecursive来自 src/core/concurrency/scheduler.recursive.js,其实现为:
schedulerProto.scheduleRecursive = function (state, action) { return this.schedule([state, action], invokeRecImmediate); };它把"状态 + 动作"打包成一次调度,由invokeRecImmediate在动作内部循环调用recurse,从而做到不爆栈的递归式循环发射。这也是RepeatSink不需要显式持有订阅对象却能响应退订的原因:调度器返回的 disposable 就是整个循环的取消句柄。
此外,模块化版本 src/modular/observable/repeat.js 给出了另一种等价实现:它把value包装成一个符合@@iterator协议、带remaining计数器的迭代器(repeatValue),再用ConcatObservable+scheduleRecursive逐个消费迭代器项并转发给下游观察者。这种"迭代器 + 串联"的结构与observableProto.repeat的实现思路一脉相承,可以对照阅读。
七、单元测试验证:行为契约一览
完整的测试用例位于 tests/observable/repeat.js,覆盖了静态工厂与原型方法的全部关键分支,是理解行为契约的最佳参考:
| 测试组 | 验证点 |
|---|---|
repeat value count zero / one / ten | 有限次数下发射个数精确,结束后onCompleted |
repeat value count dispose | 有限次数未跑完时退订立即中断 |
repeat value | repeatCount = -1即无限重复,靠退订停止 |
Repeat Observable basic / infinite | 原型方法对冷序列重放;源不完成则只订阅一次 |
Repeat Observable error | 源序列出错时,错误透传、不再重试 |
Repeat Observable throws | 观察者回调抛异常、源抛异常时正确上抛 |
Repeat Observable repeat count Basic / dispose / infinite / error / throws | 带次数的原型方法全部边界 |
例如'repeat value count ten'断言了从虚拟时刻 201 到 210 连续 10 个onNext(20x, 42),再于 210 时刻onCompleted;'Repeat Observable basic'则用xs.subscriptions.assertEqual精确断言了冷序列被反复订阅的时间窗(subscribe(200, 450)、subscribe(450, 700)……),直观展示了"序列重复 = 重新订阅"的本质。同目录下还有模块化测试 src/modular/test/repeat.js 与性能测试 tests/perf/operators/repeat.js 可供对照。
八、使用建议与注意事项
- 明确区分"重复值"与"重复序列":需要单值重放用
Rx.Observable.repeat(value, n),需要整段序列循环用source.repeat(n); - 无限重复务必可控:不传
repeatCount的序列永不onCompleted,一定要保留订阅句柄并在适当时候dispose(),否则会造成持续输出甚至资源占用; - 调度器按需显式指定:默认的同步调度器意味着
repeat会在调用线程上立即"倾泻"全部消息;若要在定时器、队列或虚拟时间(测试)环境下运行,请显式传入scheduler,例如Rx.Observable.repeat(42, 3, Rx.Scheduler.default); - 与重试机制的亲缘关系:
retry在底层就是repeat与catchError的组合,理解repeat的枚举-串联实现,等于同时理解了重试操作符的骨架。
九、配套资源定位
- API 文档:doc/api/core/operators/repeat.md
- 性能优化实现:src/core/perf/operators/repeat.js
- 原型方法实现:src/core/linq/observable/repeatproto.js
- 可枚举序列(
RepeatEnumerable):src/core/enumerable.js - 递归调度器:src/core/concurrency/scheduler.recursive.js
- 模块化实现:src/modular/observable/repeat.js
- 单元测试:tests/observable/repeat.js、src/modular/test/repeat.js、tests/perf/operators/repeat.js
- 相关操作符:
retry(src/core/linq/observable/retry.js)
【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考