RxJS中如何为pipe内的Observable实现带延迟的重试功能
这个需求完全可以实现,RxJS 内置的 retry 操作符就能直接满足你的要求,不需要自定义工具函数。
RxJS 7.0+ 版本实现(当前主流版本)
import { retry } from 'rxjs/operators'; // 提前引入操作符 notificationsWsSubject.pipe( filter((socket): socket is Socket => !!socket), switchMap(socket => fromEvent<Socket.DisconnectReason>(socket, 'disconnect')), tap(() => wsConnectedSubject.next(false)), filter(reason => (['ping timeout', 'transport close', 'transport error'] as Socket.DisconnectReason[]).includes(reason)), switchMap(() => signedInObservable), switchMap(user => forkJoin([ of(user), // 给凭证请求单独加重试逻辑 from(getNotificationsWebsocketTicket()).pipe( retry({ count: 3, // 最多重试3次 delay: 5000, // 每次重试间隔5秒 // 可选配置:只对特定类型错误重试,不需要可删除 // retryIf: (error) => error.status >= 500 }) ) ])), ).subscribe(values => { // Connect with websocket }, error => { // 只有3次重试全部失败后才会走到这里 // Throw error to user })
以上写法完全匹配你期望的 delayedRetry(3, 5000) 效果,且重试逻辑被封装在凭证请求内部,不会触发前面的socket断开监听、用户信息获取等上游逻辑,符合你的业务要求。
RxJS 6.x 及更早版本兼容实现
低版本 retry 不支持delay配置,可以用 retryWhen 组合实现相同效果:
import { retryWhen, take, concatMap, timer, throwError } from 'rxjs/operators'; // 其他逻辑不变,仅修改凭证请求部分 from(getNotificationsWebsocketTicket()).pipe( retryWhen(errors => errors.pipe( concatMap((error, retryCount) => { if (retryCount >= 3) { // 重试3次失败,抛出错误 return throwError(() => error); } // 等待5秒后重试 return timer(5000); }) )) )
内容的提问来源于stack exchange,提问作者nil
相关产品推荐
相关产品推荐

