You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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方法。

原理说明

  1. 重试触发逻辑:当map处理到'Hello'抛出错误时,retry会拦截这个错误,等待delay回调返回的Observable(这里是timer(5000))发出值后,重新订阅源Observable。
  2. 无限重试机制:retry默认不限制重试次数(除非配置count参数),因此每次遇到错误都会触发延迟重试,从而无限循环输出1、2、3。
  3. 与retryWhen的等价性:原retryWhen通过监听错误流、对每个错误延迟后触发重试;新retry的delay回调针对每个错误返回延迟Observable,最终实现的效果完全一致。

内容的提问来源于stack exchange,提问作者sandedwich

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 07:15:27