如何将含缓存逻辑的async/await函数转换为RxJS实现并确定parseApiData调用位置
Hey there! Let's work through converting your async/await cache-fetch logic to RxJS, and clear up where to place that parseApiData() method. First, let's fix a tiny bug in your original code that's easy to miss:
Your current async/await function declares it returns
Promise<Model>, butids.map(...)gives you an array ofPromise<Model>. You need to wrap that map inawait Promise.all(...)to get a singlePromise<Model[]>—that's a common gotcha with async maps!
RxJS Implementation (No subscribe() Inside the Function)
We want to return an Observable<Model[]> that emits the full list of cached/fetched Model instances once all operations are done. Here's the idiomatic way to do this with RxJS:
- For each ID, create an observable that checks the cache first. If the item exists, emit it immediately. If not, fetch from the API, parse the response, store the parsed
Modelin cache, then emit it. - Combine all these individual observables into one that emits the final array (using
forkJoin, RxJS's equivalent ofPromise.all).
Full Code:
import { forkJoin, of, Observable } from 'rxjs'; import { map, switchMap } from 'rxjs/operators'; public getListByIds(ids: number[]): Observable<Model[]> { // Create an observable stream for each ID const modelStreams = ids.map(id => { const cachedModel = this.cacheService.get(id); if (cachedModel) { // Emit the cached Model right away return of(cachedModel); } else { // Fetch, parse, cache, then emit the Model return this.http.get(`${baseUrl}/get/${id}`).pipe( // Parse raw API JSON to Model instance HERE map(rawApiData => this.parseApiData(rawApiData)), switchMap(parsedModel => { // Store the processed Model in cache this.cacheService.set(id, parsedModel); // Emit the parsed Model to the stream return of(parsedModel); }) ); } }); // Combine all streams into one that emits the full array of Models return forkJoin(modelStreams); }
Where to Call parseApiData()?
As shown in the code above, you should call parseApiData() immediately after receiving the API response (using the map operator). This is the perfect spot because:
- It converts the raw JSON from your API into a proper
Modelinstance before storing it in cache (matching your requirement to cache processedModels) - The cache ends up holding the exact type of data you want to reuse later
- The final emitted value from the observable is the processed
Model, just like your async/await code
Quick Notes:
forkJoinruns all API requests in parallel (just likePromise.all). If you need to fetch items sequentially (one after another), replaceforkJoin(modelStreams)withconcat(...modelStreams).pipe(toArray()).- By returning the observable directly, we let the caller handle subscribing (and error handling) wherever it makes sense in their code—this keeps your function clean and follows RxJS best practices.
Corrected Async/Await Version (For Reference)
Just to make sure we're on the same page, here's the fixed async/await code (including the missing cache set call):
public async getListByIds(ids: number[]): Promise<Model[]> { return Promise.all(ids.map(async id => { let cached = this.cacheService.get(id); if (!cached) { const apiData = await this.http.get(`${baseUrl}/get/${id}`).toPromise(); cached = this.parseApiData(apiData); this.cacheService.set(id, cached); // Don't forget to save the parsed Model! } return cached; })); }
内容的提问来源于stack exchange,提问作者koral

