如何以自定义交替有序方式合并两个Observable?
嘿,这个问题我之前处理有序合并数据流的时候也碰到过!你的思路其实方向是对的——用zipWith来配对两个Observable的元素,但确实差一步展平操作,才能把嵌套的Observable<Observable<Integer>>转换成我们需要的Observable<Integer>。
方法一:用flatMap展平嵌套Observable
flatMap是RxJava里处理嵌套Observable的常用操作符,它会把每个内部Observable发射的元素“平铺”到外层数据流中。针对你的场景,代码可以这么写:
Observable<Integer> obs1 = Observable.just(1, 3, 5, 7, 9); Observable<Integer> obs2 = Observable.just(2, 4, 6, 8, 10); Observable<Integer> mergedObs = obs1.zipWith(obs2, (a, b) -> Observable.just(a, b)) .flatMap(innerObs -> innerObs);
这里zipWith会把obs1和obs2的元素一一配对,每次返回一个包含两个元素的小Observable;随后flatMap会把这些小Observable的元素依次发射出来,最终得到你想要的[1,2,3,4,...10]有序数据流。
方法二:用flatMapIterable简化操作
如果你觉得返回Observable有点繁琐,还可以换个思路:在zipWith里把配对的元素放到一个集合里,再用flatMapIterable把集合里的元素逐个发射出来,代码更简洁:
Observable<Integer> mergedObs = obs1.zipWith(obs2, (a, b) -> Arrays.asList(a, b)) .flatMapIterable(list -> list);
这个写法和第一种效果完全一致,只是用集合替代了内部Observable,对于这种固定数量的元素配对场景,可读性会更好一点。
额外说明
如果你的场景对顺序要求极高(比如这个交替合并的场景),concatMap也是个不错的选择——它和flatMap的区别是会严格按顺序处理每个内部Observable/集合,不会并行处理。不过在你的这个场景里,因为每个内部的元素都是同步发射的,flatMap和concatMap的执行结果没有差异,你可以根据习惯选择。
内容的提问来源于stack exchange,提问作者Alex Kokorin

