RxJS 6.4如何实现observable队列失败后从故障位置重试3次
实现思路
RxJS 6.4 版本可以通过 concat + retry 组合实现需求,核心逻辑如下:
- 用
concat操作符保证 observable 按顺序执行,只有前一个 observable 正常完成后才会订阅下一个 - 给每个单独的 observable 加上
retry重试逻辑,仅对当前失败的任务重试,不会重置前面已经完成的任务
完整代码示例
首先导入依赖:
import { concat, of, throwError } from 'rxjs'; import { retry, tap, delay, take } from 'rxjs/operators';
基础实现(失败立即重试):
// 你的待执行 observable 数组,要求为冷 observable(订阅才触发执行) const taskList = [ of('任务1执行成功').pipe(tap(() => console.log('正在执行任务1'))), // 模拟会执行失败的任务 throwError('任务2执行出错').pipe(tap(() => console.log('正在执行任务2'))), of('任务3执行成功').pipe(tap(() => console.log('正在执行任务3'))) ]; // 给每个任务加上最多3次重试逻辑 const taskListWithRetry = taskList.map(task => task.pipe( // 重试参数3指最多额外重试3次,加上首次执行总共最多执行4次 retry(3) ) ); // 按顺序执行所有任务 concat(...taskListWithRetry).subscribe({ next: result => console.log('任务返回结果:', result), error: err => console.error('任务重试3次后仍失败,流程终止:', err), complete: () => console.log('所有任务全部执行完成') });
如果需要重试前加延迟(比如间隔1秒再重试),可以把 retry(3) 替换为 retryWhen 实现:
task.pipe( retryWhen(errors => errors.pipe( // 每次失败等待1秒再重试 delay(1000), // 最多重试3次,超过3次就抛出错误终止 take(3) ) ) )
注意事项
- 传入的 observable 必须为冷 observable,否则重试时不会重新执行任务逻辑,无法达到重试效果
- 如果单个 observable 会 emit 多个值,
concat会等该 observable 所有值发送完毕、并且触发 complete 回调后才会执行下一个任务,符合“执行完成再走下一个”的要求
内容的提问来源于stack exchange,提问作者Alexandru-Marian Constantin
相关产品推荐
相关产品推荐

