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

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 times check() actually failed.
  • The takeWhile(val => !val) will stop when it gets a true, 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 use firstValueFrom instead, 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:

  1. Run check() every 1 second
  2. Resolve the Promise with true as soon as check() succeeds, and stop all further checks
  3. 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 your check() Promise into an Observable so RxJS can handle it, and maps any errors from check() to a "failure" state.
  • scan: Keeps track of how many times we've failed consecutively. If check() 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. The true second 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 successful true emission or an error, which firstValueFrom converts 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 04:10:46