RxJS repeat操作符的实现原理是什么?如何重复Observable流?
拆解RxJS
retry操作符的核心工作原理 嘿,我来帮你捋清楚retry操作符的本质——源码看起来复杂主要是因为耦合了RxJS内部的订阅管理、工具函数,但核心逻辑其实非常直白:它的工作方式不是“重复发射已有数据”,而是当源Observable抛出错误时,重新订阅整个源流。
核心思路拆解
retry的实现可以浓缩成三个关键步骤:
- 维护重试计数器:跟踪已经进行了多少次重试,和你指定的最大次数对比。
- 封装订阅逻辑:把“订阅源Observable”这个动作做成可重复调用的函数,方便出错时再次触发。
- 错误拦截与重试判断:捕获源流的错误,根据剩余重试次数决定是重新订阅,还是把错误传递给下游。
用伪代码模拟核心逻辑
下面是去掉RxJS内部复杂封装后的核心实现,你一看就懂:
function retry(maxRetries) { return function(sourceObservable) { return new Observable((observer) => { let currentRetry = 0; // 把订阅源的逻辑封装成函数,方便重复调用 const subscribeToSource = () => { const subscription = sourceObservable.subscribe({ next: (val) => observer.next(val), // 正常数据直接传给下游 error: (err) => { if (currentRetry < maxRetries) { // 还有重试次数,计数器+1,重新订阅源 currentRetry++; subscribeToSource(); } else { // 重试耗尽,把错误传递给订阅者 observer.error(err); } }, complete: () => observer.complete() // 源流完成,下游也完成 }); return subscription; }; // 第一次订阅源流 return subscribeToSource(); }); }; }
结合你的示例看执行流程
你的示例代码是这样的:
Rx.Observable.of(2, 4, 5, 8, 10) .map(num => { if(num % 2 !== 0) { throw new Error('Odd number'); } return num; }) .retry(3) .subscribe( num => console.log(num), err => console.error(err) );
它的执行过程是:
- 第一次订阅:源流发射
2、4,到5时map抛出错误,此时currentRetry是0(小于3),计数器加1,重新订阅源流。 - 第二次订阅:再次发射
2、4,到5又出错,currentRetry变成1,继续重新订阅。 - 第三次订阅:重复上述流程,
currentRetry变成2,继续重新订阅。 - 第四次订阅:还是到
5出错,此时currentRetry已经到3(等于maxRetries),于是把错误传递给subscribe的错误回调。
最终你会在控制台看到2、4重复打印4次(1次原始订阅+3次重试),最后输出Error: Odd number。
额外注意点
retry针对的是整个源流的重新订阅,如果源是冷Observable(比如of、from),每次订阅都会从头生成数据;如果是热Observable(比如DOM点击事件),重试只是重新监听事件,不会重播之前的事件。- 如果调用
retry()不传入参数,会无限重试,直到源流成功完成。 - 源码里还处理了取消订阅的逻辑(比如下游取消订阅时,要同时取消当前的源订阅),但这属于框架的通用保障,不影响核心逻辑的理解。
内容的提问来源于stack exchange,提问作者amedina
相关产品推荐
相关产品推荐

