RxJS中retry算子能否像retryWhen一样捕获错误并实现无限重试?
用RxJS的retry算子替代废弃的retryWhen实现无限重试逻辑
问题背景
基于Lamis Chebbi所著《Reactive Patterns with RxJS for Angular》第五章「错误处理」的「重试策略」小节,原示例使用即将废弃的retryWhen算子实现了无限重试逻辑。现在需要用RxJS推荐的retry算子替代,复现相同的输出效果。
原代码实现(使用retryWhen)
服务文件 observable.service.ts
export class ObservableService { observable$ = from(['1', '2', '3', 'Hello', '100']); }
消费组件 app.component.ts
ngOnInit() { this.observableService.observable$.pipe( map((value) => { if (isNaN(value as any)) { throw new Error; } else return parseInt(value); }), retryWhen((errors) => { return errors.pipe(delayWhen(() => timer(5000))); }), ).subscribe({ next: (value) => console.log('Value emitted', value), error: (error) => console.log('Error: ', error), complete: () => console.log('Stream completed'), }); }
原输出效果
每5秒无限重复输出:
Value emitted 1 Value emitted 2 Value emitted 3
原逻辑理解
retryWhen捕获map抛出的错误,将错误作为一个通知流处理;当错误流发出通知时,会重新订阅源Observable。由于源中的'Hello'总会触发map抛出错误,因此形成了无限循环的重试逻辑。
尝试用retry算子的错误实现
尝试的组件代码
ngOnInit() { this.observableService.observable$.pipe( map((value) => { if (isNaN(value as any)) { throw new Error; } else return parseInt(value); }), retry({ delay: (error) => { return error.pipe(delayWhen(() => timer(5000))); } }), ).subscribe({ next: (value) => console.log('Value emitted', value), error: (error) => console.log('Coming from observer error handling function: ', error), complete: () => console.log('Stream completed'), }); }
报错输出
Value emitted 1 Value emitted 2 Value emitted 3 Coming from observer error handling function: TypeError: error.pipe is not a function at delay (app.component.ts:22:48) at retry.js:42:99 at OperatorSubscriber._error (OperatorSubscriber.js:23:21) at OperatorSubscriber.error (Subscriber.js:40:18) at OperatorSubscriber._next (OperatorSubscriber.js:16:33) at OperatorSubscriber.next (Subscriber.js:31:18) at Observable._subscribe (innerFrom.js:51:24) at Observable._trySubscribe (Observable.js:37:25) at Observable.js:31:30 at errorContext (errorContext.js:19:9)
正确的retry算子实现及原理
修复后的组件代码
import { timer } from 'rxjs'; // ...其他导入 ngOnInit() { this.observableService.observable$.pipe( map((value) => { // 用Number()更准确判断是否为数字 if (isNaN(Number(value))) { throw new Error('Invalid numeric value'); } // 加上基数参数避免意外行为 return parseInt(value, 10); }), retry({ // delay回调返回一个Observable,发出后触发重试 delay: () => timer(5000) }), ).subscribe({ next: (value) => console.log('Value emitted', value), error: (error) => console.log('Coming from observer error handling function: ', error), complete: () => console.log('Stream completed'), }); }
错误原因解析
你之前的代码错误在于:retry算子的delay回调参数是单个错误对象,而非retryWhen接收的错误Observable流。因此error.pipe()会报错,因为普通Error对象没有pipe方法。
原理说明
- 重试触发逻辑:当
map处理到'Hello'抛出错误时,retry会拦截这个错误,等待delay回调返回的Observable(这里是timer(5000))发出值后,重新订阅源Observable。 - 无限重试机制:
retry默认不限制重试次数(除非配置count参数),因此每次遇到错误都会触发延迟重试,从而无限循环输出1、2、3。 - 与
retryWhen的等价性:原retryWhen通过监听错误流、对每个错误延迟后触发重试;新retry的delay回调针对每个错误返回延迟Observable,最终实现的效果完全一致。
内容的提问来源于stack exchange,提问作者sandedwich
相关产品推荐
相关产品推荐

