RxJS中forkJoin合并BehaviorSubject与Observable无响应解决方案
forkJoin 只会在所有传入的Observable全部执行complete之后,才会发送每个流的最后一个值。你代码里的usersChange是基于BehaviorSubject转换的长期存活流,只要不手动触发complete,这个流永远不会结束,forkJoin会一直等待流完成,自然不会发出任何值。
另外你贴的代码还有个书写错误:forkJoin执行后的返回Observable没有赋值给this.getUsersDataAndLoggedUserSubscription,你实际订阅的是未定义的变量,这也会导致逻辑不执行。
根据你的业务场景选对应操作符即可:
场景1:仅首次获取两个流的当前值,后续用户列表更新不触发
如果只需要初始化时拿一次当前用户列表数据+当前登录用户信息,不需要响应后续的用户列表更新,用zip搭配take(1)实现:
import { zip, take } from 'rxjs'; // 记得把合并后的Observable赋值给订阅变量,方便后续销毁取消订阅 this.getUsersDataAndLoggedUserSubscription = zip( // take(1)表示只取BehaviorSubject当前存的最新值,取完后自动结束当前分支流 this.saveUsersSubject.usersChange.pipe(take(1)), this.api.getCurrentUser() ).subscribe(([userData, loggedUser]) => { this.actionedUsers = userData; this.currentUser = loggedUser; });
这个方案的效果和你原本想用forkJoin达到的一次性取值效果完全一致。
场景2:后续用户列表更新时,自动用最新列表+已拿到的用户信息触发回调
如果希望每次调用updatedUsers推送新的用户列表时,订阅逻辑能自动拿到最新列表执行更新,不需要重复请求当前用户接口,用combineLatest是最优解:
import { combineLatest } from 'rxjs'; this.getUsersDataAndLoggedUserSubscription = combineLatest([ this.saveUsersSubject.usersChange, this.api.getCurrentUser() ]).subscribe(([userData, loggedUser]) => { this.actionedUsers = userData; this.currentUser = loggedUser; });
combineLatest的执行逻辑是:等所有传入的流都至少发出过1个值之后,任意一个流推送新值,都会把所有流的最新值组合后发给订阅者。你的BehaviorSubject本身就有初始空数组作为第一个值,只要getCurrentUser接口返回,就会立刻触发第一次回调,后续用户列表更新时会自动触发后续回调,不会重复发起用户信息请求。
注意:不管用哪种方案,记得在组件/服务的销毁周期里调用
this.getUsersDataAndLoggedUserSubscription.unsubscribe()取消订阅,避免内存泄漏。
内容的提问来源于stack exchange,提问作者Bojidari

