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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 19:15:04