如何在RxJS中重置retryWhen的递增重试等待时间间隔?
解决RxJS retryWhen递增重试间隔的重置问题
你遇到的问题核心在于:当前代码里的Observable.range(1,4)是一次性订阅的,重试计数会一直累加,不会在连接成功后重置。要实现“连接成功后下次出错重新从1秒开始重试”的效果,我们需要让每次错误重试周期都重新初始化计数。
方案一:利用switchMap+takeUntil重置重试序列
这个思路是:每次源发出错误时,创建一个全新的从1开始的延迟序列;当源成功连接(发出正常数据)时,终止当前的延迟序列,确保下次出错时重新生成序列。
首先要注意共享源Observable,避免重复订阅带来的不必要开销:
// 先共享源,确保多个订阅复用同一个流 const sharedSocket = socketResponse.share(); sharedSocket.retryWhen(attempts => { return attempts.switchMap(() => { // 每次错误触发时,创建新的递增延迟序列(1-4秒) return Observable.range(1, 4).concatMap(i => { console.log(`[socket] Wait ${i} seconds, then retry!`); if (i === 4) { console.log(`[socket] maxReconnectAttempts ${i} reached!`); // 可选:达到最大重试次数后抛出错误,终止重试 // return Observable.throw(new Error('Max reconnect attempts exceeded')); } return Observable.timer(i * 1000); }).takeUntil(sharedSocket); // 源成功时终止当前延迟序列,重置计数 }); });
方案二:用scan跟踪重试计数并重置
这个方法通过scan操作符维护当前重试次数,当源成功发出数据时重置计数,错误发生时递增计数,逻辑更清晰可控:
const sharedSocket = socketResponse.share(); sharedSocket.retryWhen(attempts => { // 合并源的成功事件和错误事件,统一处理计数 return sharedSocket .merge(attempts.map(error => ({ type: 'error', error }))) .scan((currentCount, event) => { if (event.type === 'error') { // 错误发生,计数+1 return currentCount + 1; } else { // 源成功,重置计数为0 return 0; } }, 0) .filter(count => count > 0) // 只处理错误触发的计数 .mergeMap(count => { if (count > 4) { console.log(`[socket] maxReconnectAttempts 4 reached!`); return Observable.throw(new Error('Max retries reached')); } console.log(`[socket] Wait ${count} seconds, then retry!`); return Observable.timer(count * 1000); }); });
效果验证
- 首次连接错误:
[socket] Wait 1 seconds, then retry! [socket] Wait 2 seconds, then retry! - 连接成功后再次出错:
[socket] Wait 1 seconds, then retry! [socket] Wait 2 seconds, then retry!
完美符合你的期望,每次成功后都会重置重试间隔。
内容的提问来源于stack exchange,提问作者durgesh rao
相关产品推荐
相关产品推荐

