如何用RxJS实现带超时的异步请求流及特定Observable流程
RxJS Solutions for Timeout Handling & Custom Observable Flow
Let's break down your requirements into two clear, practical solutions that align with your async/await logic.
Part 1: Timeout Handling for concat/merge Async Requests
When working with concat or merge in RxJS, you can apply timeouts in two useful ways: per-request timeouts (each async call has its own deadline) or a global timeout for the entire sequence. Here's how to implement both:
Per-Request Timeout
Use the timeout operator on each individual Observable before passing them to concat/merge, plus catchError to handle timeout errors gracefully:
import { concat, merge, from, throwError } from 'rxjs'; import { timeout, catchError } from 'rxjs/operators'; // Sample async requests (replace with your actual API calls) const fetchUser = () => fetch('/api/user').then(res => res.json()); const fetchPosts = () => fetch('/api/posts').then(res => res.json()); // Concat with per-request timeouts concat( from(fetchUser()).pipe( timeout(2000), // 2-second timeout for user request catchError(() => throwError(() => new Error('User fetch timed out'))) ), from(fetchPosts()).pipe( timeout(4000), // 4-second timeout for posts request catchError(() => throwError(() => new Error('Posts fetch timed out'))) ) ).subscribe({ next: data => console.log('Received:', data), error: err => console.error('Failure:', err.message) }); // Merge works the same way—just replace concat with merge merge( from(fetchUser()).pipe(timeout(2000)), from(fetchPosts()).pipe(timeout(4000)) ).subscribe(...);
Global Sequence Timeout
If you want a single timeout for the entire concat/merge sequence (e.g., both requests combined must finish within 10 seconds), apply timeout to the final stream:
concat(from(fetchUser()), from(fetchPosts())) .pipe(timeout(10000)) // Total 10-second limit for both requests .subscribe({ next: data => console.log('Received:', data), error: () => console.error('Total sequence timed out') });
Part 2: Custom Observable Matching Your Async/Await Flow
Your async/await code involves fetching game info, filtering out invalid results, waiting a calculated delay, then running another async call with timeout. Here's how to replicate this in RxJS, including converting to Promise where required:
Full RxJS Stream Implementation
This uses operators to chain the entire flow seamlessly without breaking into Promises unless necessary:
import { from, EMPTY } from 'rxjs'; import { timeout, filter, delayWhen, switchMap, tap, catchError } from 'rxjs/operators'; // Assume these are your existing functions const function1 = async () => { /* Returns undefined or { start: Date, _id: string } */ }; const function3 = async (id: string) => { /* Your async logic here */ }; const getDelayTime = (start: Date) => { /* Calculate delay in ms */ }; // Build the Observable stream from(function1()).pipe( // Add timeout to function1 call timeout(3000), // Filter out falsy gameInfo (stop stream if undefined) filter(gameInfo => !!gameInfo), // Wait for the calculated delay before proceeding delayWhen(gameInfo => from(Promise.resolve()).pipe(delay(getDelayTime(gameInfo.start)))), // Run function3 with timeout, log when done switchMap(gameInfo => from(function3(gameInfo._id)).pipe( timeout(5000), tap(() => console.log('Done')) )), // Handle any errors (timeouts or failed requests) catchError(err => { console.error('Error:', err); return EMPTY; }) ).subscribe();
With Promise Conversion (As Per Requirement)
If you need to convert part of the flow to a Promise (e.g., getting gameInfo as a Promise first), use firstValueFrom to convert the Observable, then proceed with the delay and function3 call:
import { from, firstValueFrom } from 'rxjs'; import { timeout } from 'rxjs/operators'; // Get gameInfo as a Promise with timeout firstValueFrom(from(function1()).pipe(timeout(3000))) .then(gameInfo => { if (!gameInfo) return; // Filter out invalid data // Wait for calculated delay, then run function3 with timeout setTimeout(() => { from(function3(gameInfo._id)).pipe(timeout(5000)) .subscribe({ next: () => console.log('Done'), error: err => console.error('Function3 failed:', err) }); }, getDelayTime(gameInfo.start)); }) .catch(err => console.error('Function1 timed out or failed:', err));
This second approach closely mirrors your original async/await structure while adding the required timeout handling at each step.
内容的提问来源于stack exchange,提问作者Grynets

