如何合并flatMap的发射值?switchMap后使用flatMap的技术咨询
先回答你第一个问题:flatMap(RxJS中更推荐叫mergeMap,两者是同一个操作符)是怎么合并发射值的?
简单来说,flatMap会把上游每一个发射出来的值,都转换成一个新的Observable,然后同时订阅所有这些新的Observable。不管这些Observable什么时候返回数据,只要有值产生,就立刻把值推给下游,完全按照数据产生的时间顺序合并,不会等前一个Observable完成才处理下一个。举个例子:如果上游先发射idA,flatMap转换成一个需要2秒返回的Observable;接着上游发射idB,flatMap转换成一个需要1秒返回的Observable。那下游会先收到idB对应的结果,再收到idA对应的结果。
接下来看你的代码场景问题:
从你给出的代码片段来看,核心问题出在switchMap的回调返回值上。switchMap要求你返回一个单个Observable,但你现在返回的是ids.map(id => this.db.object(...))——这是一个Observable数组,RxJS没法正确处理这种类型,这会导致整个数据流出现异常,甚至可能影响上游map操作的执行(不过你说id的日志被调用了1次,说明上游的map其实已经执行过,只是可能你没订阅整个流,导致第一个日志没被看到?先解决代码的核心问题:
修正代码的两种方案
根据你的需求,你可以选择以下两种方式来合并多个对象Observable:
方案1:等待所有对象数据都返回后,一次性输出结果数组
如果需要等所有id对应的对象都获取完成,再把结果作为一个数组发射给下游,用forkJoin:
import { forkJoin, Observable } from 'rxjs'; import { map, switchMap } from 'rxjs/operators'; // ... return this.db.list(`UserPlaces/${this.authData.auth.auth.currentUser.uid}`, { query: { orderByChild: 'deleted', equalTo: false } }).pipe( map((locations: any) => { console.log(JSON.stringify(locations,null,2)); // 现在应该会触发 return locations.map(location => location.$key); }), switchMap(ids => { // 处理空数组的情况,避免forkJoin没有输入报错 if (ids.length === 0) { return Observable.of([]); } // 把每个id转换成对应的对象Observable const objectObservables = ids.map(id => this.db.object(`你的对象路径/${id}`).valueChanges() // 假设用valueChanges获取数据 ); // forkJoin会等所有Observable完成后,返回结果数组 return forkJoin(objectObservables); }) );
方案2:每个对象数据一返回就立刻发射,不需要等全部完成
如果不需要等待所有结果,只要某个id的对象数据就绪就立刻推给下游,用merge把Observable数组合并成单个Observable:
import { merge } from 'rxjs'; import { map, switchMap } from 'rxjs/operators'; // ... return this.db.list(`UserPlaces/${this.authData.auth.auth.currentUser.uid}`, { query: { orderByChild: 'deleted', equalTo: false } }).pipe( map((locations: any) => { console.log(JSON.stringify(locations,null,2)); // 现在应该会触发 return locations.map(location => location.$key); }), switchMap(ids => { const objectObservables = ids.map(id => this.db.object(`你的对象路径/${id}`).valueChanges() ); // merge会订阅所有Observable,按数据产生顺序发射每个结果 return merge(...objectObservables); }) );
关于第一个日志没被调用的补充说明
你提到console.log(JSON.stringify(locations,null,2))没被调用,但id的日志被调用了1次,这有点矛盾。大概率是因为整个Observable没有被订阅——RxJS的Observable是“冷”的,只有当你调用.subscribe()方法时,才会触发整个数据流的执行。比如你在组件里调用这个方法时,需要加上:
this.yourService.getYourData().subscribe( result => console.log('最终结果', result), error => console.error('出错了', error) );
如果没有订阅,上游的map操作不会执行,也就不会输出日志,建议检查一下订阅逻辑,确保整个数据流被正确订阅。
内容的提问来源于stack exchange,提问作者Khaled Ramadan

