如何使用RxJS关联映射两个异步数据流并输出到目标Subject?
RxJS 实现双异步流映射方案
首先对齐你的代码上下文和逻辑:你需要将users流发出的用户数组,结合groups$流发出的权限数据,给每个用户匹配对应的value属性,最终将处理后的数据发送到mapped$中。
根据你的业务需要,有两种常用实现方式:
方案1:任意流更新都重新计算映射结果
如果你的需求是只要users或者groups$任意一个流发出新值,就重新计算最新的映射结果,使用combineLatest操作符即可,它会等待所有输入流至少发出过一次值后,每次任意流更新都把所有流的最新值组合下发:
import { combineLatest } from 'rxjs'; import { map } from 'rxjs/operators'; // 组合两个流的最新值 combineLatest([users, groups$]).pipe( map(([userList, groupList]: [User[], Rights[]]) => { // 建议返回新对象避免修改原数据产生副作用,若需要直接修改原对象可改用forEach写法 return userList.map(user => ({ ...user, value: groupList.find(g => g.user === user.id)?.value })) }) // 直接订阅后将结果推送到mapped$ ).subscribe(mapped$);
如果需要在groups$还没发出初始值时就可以正常运行,可以给groups$加上默认值:
import { startWith } from 'rxjs/operators'; combineLatest([ users, // 默认值可以根据业务调整,这里默认是空数组 groups$.pipe(startWith([] as Rights[])) ]).pipe( // 同上map逻辑 ).subscribe(mapped$)
方案2:仅当用户流更新时才计算映射结果
如果你的需求是只有users流发出新值的时候,才取groups$的当前最新值做映射,groups$自身更新不需要触发重计算,可以使用withLatestFrom操作符:
import { withLatestFrom, map } from 'rxjs/operators'; users.pipe( withLatestFrom(groups$), map(([userList, groupList]: [User[], Rights[]]) => { return userList.map(user => ({ ...user, value: groupList.find(g => g.user === user.id)?.value })) }) ).subscribe(mapped$);
注意事项
- 原命令式代码直接修改了原
user对象,容易产生副作用,示例中采用返回新对象的写法更符合RxJS的函数式编程规范,如果你确实需要修改原对象,可以把map替换为forEach修改后再返回原数组即可。 - 上述代码适用于RxJS 6及以上版本。
内容的提问来源于stack exchange,提问作者user15361861
相关产品推荐
相关产品推荐

