RxJS 4 的 `Rx.Observable.repeat` 深入解析:固定值重复发射、调度器与源码实现
2026/9/21 16:22:59 网站建设 项目流程

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 次再串联起来。

二、参数详解与默认值

参数类型必填默认值说明
valueAny要重复发射的元素,可以是任意 JS 值(数字、字符串、对象等)
repeatCountNumber-1重复发射的次数;不传或传null时等价于无限重复
schedulerSchedulerScheduler.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.RepeatSinkrun方法(发射循环本体)

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 valuerepeatCount = -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 可供对照。

八、使用建议与注意事项

  1. 明确区分"重复值"与"重复序列":需要单值重放用Rx.Observable.repeat(value, n),需要整段序列循环用source.repeat(n)
  2. 无限重复务必可控:不传repeatCount的序列永不onCompleted,一定要保留订阅句柄并在适当时候dispose(),否则会造成持续输出甚至资源占用;
  3. 调度器按需显式指定:默认的同步调度器意味着repeat会在调用线程上立即"倾泻"全部消息;若要在定时器、队列或虚拟时间(测试)环境下运行,请显式传入scheduler,例如Rx.Observable.repeat(42, 3, Rx.Scheduler.default)
  4. 与重试机制的亲缘关系retry在底层就是repeatcatchError的组合,理解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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询