如何重试RxJS中的内部Observable?附实现代码
Hey there, let's break down your current RxJS implementation and figure out where it might fall short, plus how to make it better.
Issues with the Current Implementation
Recursive
catchleads to messy logic & potential memory bloat
Yoursmfunction usescatchto recursively call itself when an error occurs. While this technically triggers retries, it creates nested Observables every time a retry happens—making the code hard to debug and risking memory leaks if retries pile up. Worse, there's no built-in limit to retries, so this could loop infinitely ifsm$keeps throwing errors.Unintended cancellation from outer
switchMap
The outerinterval(500)usesswitchMap(sm). Since your retry delay is 1000ms (longer than the interval), a new interval event will cancel the ongoing retry process for the previous value. That means you might never get a successful emission for a value if retries take longer than 500ms, which is probably not what you want.Unclear error handling scope
Thecatchis applied after the firstswitchMap(sm$), but the recursive call re-runs the entireObservable.of(val).switchMap(sm$)chain. This makes it impossible to track how many times a single value has been retried, and adds unnecessary complexity to the stream.
Optimized Implementation
Let's fix this using RxJS's built-in retryWhen operator—it's designed exactly for retry-on-error scenarios, gives you fine-grained control over delays and retry limits, and keeps the code clean. We'll also address the outer stream behavior to match your needs.
First, keep your original sm$ function (no changes needed here):
function sm$(val) { if (Math.random() > .4) { return Rx.Observable.throw(val); } else { return Rx.Observable.of(val); } }
Now, rewrite the sm function with retryWhen for controlled retries:
function sm(val) { return Rx.Observable.of(val) .switchMap(sm$) // Retry up to 3 times with 1s delay between each attempt .retryWhen(errors => errors.pipe( Rx.Observable.delay(1000), Rx.Observable.take(3), // Throw an error if we hit the retry limit Rx.Observable.concat(Rx.Observable.throw(`Max retries reached for value: ${val}`)) ) ); }
Adjust the outer stream based on your desired behavior:
- If you want to preserve ongoing retries (don't cancel them when a new interval event arrives), use
concatMapinstead ofswitchMap:Rx.Observable.interval(500) .concatMap(sm) // Waits for the previous observable to complete before starting the next .take(5) .subscribe( val => console.log("val:", val), err => console.log("Error:", err) ); - If you do want to cancel old retries when new interval values come in, keep
switchMapbut ensure your retry delay is shorter than the interval (or accept that longer retries will be canceled).
Key Improvements
- Controlled retries:
retryWhenlets you define exactly how many times to retry, how long to wait between attempts, and what happens when retries are exhausted. - Cleaner, maintainable logic: No recursive calls—everything is handled via standard RxJS operators, making the code easier to read and debug.
- Predictable outer stream behavior: Choosing
concatMapvsswitchMaplets you explicitly control whether ongoing retries are kept or canceled when new values arrive.
内容的提问来源于stack exchange,提问作者jigfox

