RxJava中combineLatest如何判断哪个Observable发射了数据?
解答RxJava combineLatest相关的两个问题
Hey there! Let's tackle your two RxJava questions step by step:
问题1:RxJava的combineLatest()如何知晓哪个Observable发射了值?
默认情况下,combineLatest()只会返回它合并的各个Observable的最新值——并不会直接告诉你是哪个Observable触发了这次合并操作。要追踪值的来源,你需要在合并前给每个Observable的发射值打上来源标记,具体操作如下:
- 先定义一个简单的枚举来区分不同的Observable:
enum ObservableSource { LEFT_DEPT, RIGHT_PERSON }
- 使用
map()操作符给每个原始Observable的发射值包裹上来源标记:
Observable<Pair<ObservableSource, Department>> taggedLeft = leftObservable .map(dept -> Pair.create(ObservableSource.LEFT_DEPT, dept)); Observable<Pair<ObservableSource, Person>> taggedRight = rightObservable .map(person -> Pair.create(ObservableSource.RIGHT_PERSON, person));
- 现在合并这些带标记的Observable,当合并发生时,你就能判断是哪个(或哪些)Observable发射了新值:
Observable.combineLatest(taggedLeft, taggedRight, (leftPair, rightPair) -> { // 获取原始数据 Department currentDept = leftPair.second; Person currentPerson = rightPair.second; // 判断来源(注意:两个Observable可能同时发射新值) boolean leftChanged = leftPair.first == ObservableSource.LEFT_DEPT; boolean rightChanged = rightPair.first == ObservableSource.RIGHT_PERSON; // 这里编写你的业务逻辑 return ...; });
问题2:如何判断是哪个Observable发生了变化?
针对你给出的现有代码,有两种实用方案可以追踪到底是哪个Observable触发了后续逻辑:
方案一:用scan()追踪上一次状态
利用scan()操作符保存上一次发射的Pair,通过对比当前值与上一次的值,判断哪个数据发生了变化:
首先创建一个辅助类,用来封装当前数据和变化标记(如果使用Java 16+,也可以用自定义record简化):
class PairWithChanges { private final Pair<Department, Person> currentPair; private final boolean deptChanged; private final boolean personChanged; // 构造方法、getter方法省略 }
然后将其整合到你的代码链中:
Observable.combineLatest(leftObservable, rightObservable, Pair::new) .scan((previous, current) -> { boolean isPersonChanged = !Objects.equals(previous.second, current.second); boolean isDeptChanged = !Objects.equals(previous.first, current.first); return new PairWithChanges(current, isDeptChanged, isPersonChanged); }) .switchMap(pairWithChanges -> { Department dept = pairWithChanges.getCurrentPair().first; Person person = pairWithChanges.getCurrentPair().second; if (pairWithChanges.isPersonChanged()) { // 当person发生变化时执行的逻辑 } return repository.loadHourlyTable(dept, person); });
方案二:提前打标记并过滤重复值
这种方式可以更精准地追踪变化,只在值真正改变时才发射,同时保留来源标记:
// 仅当department变化时才发射,并打上标记 Observable<Pair<ObservableSource, Department>> distinctLeft = leftObservable .distinctUntilChanged() .map(dept -> Pair.create(ObservableSource.LEFT_DEPT, dept)); // 仅当person变化时才发射,并打上标记 Observable<Pair<ObservableSource, Person>> distinctRight = rightObservable .distinctUntilChanged() .map(person -> Pair.create(ObservableSource.RIGHT_PERSON, person)); // 合并并追踪最新状态与变化来源 Observable.combineLatest(distinctLeft, distinctRight, (leftPair, rightPair) -> { // 你可以在这里记录最后变化的来源,或者同时传递两个标记 return new CombinedState(leftPair.second, rightPair.second, ObservableSource.RIGHT_PERSON); }) .switchMap(state -> { if (ObservableSource.RIGHT_PERSON.equals(state.getLastChangedSource())) { // 当person发生变化时执行的逻辑 } return repository.loadHourlyTable(state.getDept(), state.getPerson()); });
注意: 如果两个Observable同时发射新值,combineLatest会合并这两个新值——如果你的业务场景存在这种情况,需要针对性处理这个边缘场景。
内容的提问来源于stack exchange,提问作者Sergey Shustikov
相关产品推荐
相关产品推荐

