使用forkJoin配合AngularFire执行多并发查询无返回结果求助
问题根因
- 你的
PermitBrowserService是Angular单例服务,所有getData调用共享同一个contactName$行为主体和permitData$可观察对象实例:循环调用getData时,连续给contactName$推送新值,switchMap会直接取消前一次未完成的查询,最终只剩下最后一次查询在运行,无法拿到多个条件的查询结果。 - AngularFire返回的
snapshotChanges()是热可观察对象,会持续监听数据变动不会主动结束,而forkJoin要求所有输入的可观察对象都至少发出一次值且正常完成才会触发回调,因此永远不会返回结果。
修复方案
首先改造PermitBrowserService,移除共享全局状态,每次调用返回独立的查询可观察对象,根据是否需要实时更新选择对应实现:
import { take } from 'rxjs/operators'; import { forkJoin, combineLatest } from 'rxjs'; export class PermitBrowserService { constructor( public db: AngularFireDatabase, ) {} // 单联系人查询,兼容原有逻辑 getData(contactNameFilter?: string | null): Observable<WellPermit[]> { const queryRef = contactNameFilter ? this.db.list('/permits', ref => ref.orderByChild('contactName').equalTo(contactNameFilter)) : this.db.list('/permits'); return queryRef.snapshotChanges().pipe( // 不需要实时更新则保留take(1),让可观察对象拿到第一次结果后主动完成,适配forkJoin要求 take(1), map(changes => { return changes.map(c => { const data = c.payload.val() as WellPermit; const id = c.key; return { id, ...data }; }) }) ); } // 批量查询多联系人,仅获取一次结果 getBatchData(contactNames: string[]): Observable<WellPermit[]> { const requests = contactNames.map(name => this.getData(name)); return forkJoin(requests).pipe( // 拍平多个查询的结果数组 map(results => results.flat()) ); } // 批量查询多联系人,持续监听实时变动 getBatchRealtimeData(contactNames: string[]): Observable<WellPermit[]> { const requests = contactNames.map(name => // 移除take(1),持续监听数据变动 this.getData(name).pipe(take(Infinity)) ); // 用combineLatest替代forkJoin,不需要可观察对象完成即可返回最新结果 return combineLatest(requests).pipe( map(results => results.flat()) ); } }
调用示例:
contactNamesFilter: string[] = ['Steve', 'Brandon']; // 仅获取一次结果 this.permitBrowserService.getBatchData(contactNamesFilter) .subscribe(console.log); // 持续监听实时变动 this.permitBrowserService.getBatchRealtimeData(contactNamesFilter) .subscribe(console.log);
内容的提问来源于stack exchange,提问作者Tim Denning
相关产品推荐
相关产品推荐

