RXJava2合并Observable数据异常,寻求技术帮助
Hey there! Let's break down what's happening with your code and fix that data overwriting issue.
为什么merge没达到你想要的效果?
The merge operator's job is to combine the emission streams of multiple Observables. In your case, it's sending two separate Users objects to your Observer: first the one from your API, then the custom one you created.
Chances are your Observer's onNext method is replacing the existing data with each new Users object it receives. That's why you only see the custom data at the end—it's the last thing emitted.
正确的做法:合并两个Users的Items列表
What you actually want is to take the Items from both Users objects, combine them into a single list, then emit that as one Users object. Here are two solid ways to do this:
方法1:用zip操作符(最适合单发射Observable)
zip waits for both Observables to emit their data, then lets you combine them into a single result. Perfect for your case since both your Observables emit one Users object each.
Observable.zip(baseDataObservable, usersObservable, (apiUsers, customUsers) -> { // 创建合并后的Items列表 List<Items> combinedItems = new ArrayList<>(); // 先添加API返回的数据(如果存在) if (apiUsers.getItems() != null) { combinedItems.addAll(apiUsers.getItems()); } // 再添加自定义数据 if (customUsers.getItems() != null) { combinedItems.addAll(customUsers.getItems()); } // 创建包含合并列表的新Users对象 Users combinedUsers = new Users(); combinedUsers.setItems(combinedItems); return combinedUsers; }) .subscribeOn(Schedulers.newThread()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(observer);
方法2:用concatMap(适合需要顺序依赖的场景)
如果你希望确保API数据先加载完成,再添加自定义数据,concatMap是不错的选择——它会按顺序处理Observable:
baseDataObservable.concatMap(apiUsers -> { // 获取到API数据后,与自定义数据合并 return usersObservable.map(customUsers -> { List<Items> combinedItems = new ArrayList<>(); if (apiUsers.getItems() != null) { combinedItems.addAll(apiUsers.getItems()); } if (customUsers.getItems() != null) { combinedItems.addAll(customUsers.getItems()); } Users combinedUsers = new Users(); combinedUsers.setItems(combinedItems); return combinedUsers; }); }) .subscribeOn(Schedulers.newThread()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(observer);
额外检查:你的Observer逻辑
别忘了检查Observer的onNext实现。如果它是类似this.currentUsers = receivedUsers;这样的写法,每次收到新的Users对象都会覆盖旧数据。用上面的方案后,Observer只会收到一个包含所有合并数据的Users对象,就不会出现覆盖问题了。
内容的提问来源于stack exchange,提问作者Rohit

