Rxjs中takeWhile结合toPromise的实现问题求助
Hey there! Let's work through this RxJS problem together—your requirements are clear, and I can spot a few gaps in your current code that are keeping it from meeting all three needs. Let's break this down step by step.
First, let's look at the issues with your existing code:
- If
check()is an asynchronous function (returns a Promise), you're calling it synchronously here, which means you won't wait for its result before deciding what to return. If it's synchronous, the failure logic is tied to the interval count (i===4) instead of how many timescheck()actually failed. - The
takeWhile(val => !val)will stop when it gets atrue, but there's no proper handling for counting consecutive failed checks up to 5 times before rejecting. - Also, note that
toPromise()is deprecated in newer RxJS versions—we'll usefirstValueFrominstead, which is the recommended way to convert an Observable to a Promise.
Here's a solution that meets all your requirements:
Let's assume check() returns a Promise<boolean> (if it's synchronous, we can adjust this slightly). This code will:
- Run
check()every 1 second - Resolve the Promise with
trueas soon ascheck()succeeds, and stop all further checks - Reject the Promise with an error if
check()fails 5 times in a row
import { interval, throwError, of, from } from 'rxjs'; import { mergeMap, scan, takeWhile, first, map, catchError } from 'rxjs/operators'; // Replace this with your actual check function const check = async () => { // Example: simulate random success/failure return Math.random() > 0.7; }; const runChecks = () => { return firstValueFrom( interval(1000).pipe( // Convert check's Promise to an Observable, handle any errors from check as failures mergeMap(() => from(check()).pipe( map(success => ({ success })), // If check throws an error, treat it as a failure catchError(() => of({ success: false })) ) ), // Track the number of consecutive failures scan((acc, current) => { if (current.success) { // Reset failure count if we succeed return { failedCount: 0, success: true }; } // Increment failure count on each failure return { failedCount: acc.failedCount + 1, success: false }; }, { failedCount: 0, success: false }), // Continue checking only if we haven't succeeded and haven't failed 5 times yet takeWhile(state => !state.success && state.failedCount < 5, true), // Grab the first value that triggers our stop condition first(), // Resolve with true if successful, reject if we hit 5 failures mergeMap(state => { if (state.success) { return of(true); } return throwError(() => new Error('Check failed 5 consecutive times')); }) ) ); }; // Usage example runChecks() .then(result => console.log('Success!', result)) .catch(err => console.error('Error:', err.message));
Let's walk through what each part does:
interval(1000): Emits a value every 1 second to trigger a new check.mergeMap(() => from(check())): Converts yourcheck()Promise into an Observable so RxJS can handle it, and maps any errors fromcheck()to a "failure" state.scan: Keeps track of how many times we've failed consecutively. Ifcheck()succeeds, it resets the count; if it fails, it increments the count.takeWhile: Continues the stream only if we haven't succeeded and haven't failed 5 times. Thetruesecond argument ensures we emit the final state that triggers the stop (either success or the 5th failure).first(): Ensures we only take the first stop-triggering value (since we don't need any further emissions once we hit our condition).mergeMap: Converts the final state into either a successfultrueemission or an error, whichfirstValueFromconverts into a resolved or rejected Promise.
If check() is synchronous:
Just replace from(check()) with of(check())—the rest of the logic stays the same.
内容的提问来源于stack exchange,提问作者Lin Shen
相关产品推荐
相关产品推荐

