RxJava2 基于两个Observable构建Observable<List<Schedule>>遇阻
解决RxJava2中合并两个Observable生成List的问题
嘿,我完全懂你现在卡在哪了——想把频道列表和对应节目单整合成一个节目单列表,但toList()死活不触发onNext()对吧?你说的没错,核心问题就是toList()的特性:它必须等上游Observable完全完成(也就是调用了onComplete()),才会把收集到的所有数据打包成列表发射出来。咱们一步步调整你的数据流就能搞定:
问题根源拆解
你之前的实现大概率是没正确拆解频道列表,或者上游Observable没正常发送完成信号。比如如果直接对getChannels()返回的Observable<List<Channel>>做操作,却没把列表拆成单个Channel的流,那toList()只会等整个频道列表的完成信号,而不会等待每个频道对应节目单的完成,自然不会发射结果。
正确实现方案
我们把数据流拆成几个关键步骤,确保每个环节都能正确传递完成信号:
- 拆解频道列表:把
Observable<List<Channel>>转成逐个发射Channel对象的Observable,让每个频道都能单独处理。 - 获取单频道节目单:对每个频道调用
getSchedule(),拿到对应节目单的Observable。 - 收集所有节目单:用
toList()收集所有节目单,只要所有单个节目单的流都完成,toList()就会发射最终的列表。
代码示例
Observable<List<Schedule>> getCombinedScheduleList() { return getChannels() // 把频道列表拆成逐个发射的Channel对象 .flatMapIterable(channelList -> channelList) // 为每个频道获取节目单,要顺序处理就换成concatMap .flatMap(channel -> getSchedule(channel.getId())) // 收集所有Schedule到List,所有上游事件完成后发射 .toList() // 把Single转成Observable(如果你的场景需要Observable类型) .toObservable(); }
关键注意事项
- 确保上游流正常完成:一定要保证
getChannels()发射完频道列表后会调用onComplete(),同时每个getSchedule(channelId)在发射完对应节目单后也会调用onComplete()——只有这些完成信号都到位,toList()才会触发。 - 异常兜底处理:如果某个频道的
getSchedule()可能出错,可以加onErrorResumeNext()避免整个流中断:.flatMap(channel -> getSchedule(channel.getId()) .onErrorResumeNext(throwable -> Observable.empty())) // 出错时发射空流,不影响整体 - 顺序控制:如果需要节目单和频道的顺序完全一致,把
flatMap换成concatMap,它会按顺序处理每个频道的节目单请求,而flatMap是并行处理的。
内容的提问来源于stack exchange,提问作者w00ly
相关产品推荐
相关产品推荐

