如何实现未知数量带副作用的Observable Supplier的延迟链式flatMap调用?
解决方案
要动态处理任意数量的这类供应商,我们可以用RxJava的reduce操作符(或Java Stream的reduce)来动态构建延迟触发的链式结构,核心是把每个供应商通过flatMap逐个串联,保证只有上游Observable发射数据时才触发下游供应商的get()。
方法一:Java Stream + Reduce实现
public static <T> Observable<T> chainSuppliers(Collection<Supplier<Observable<T>>> suppliers) { if (suppliers.isEmpty()) { return Observable.empty(); } // 把多个供应商逐步合并成一个链式的供应商 return suppliers.stream() .reduce((prevSupplier, nextSupplier) -> () -> prevSupplier.get().flatMap(ignored -> nextSupplier.get()) ) .orElse(Observable::empty) .get(); }
使用示例
// 把所有供应商放进集合 List<Supplier<Observable<String>>> supplierList = Arrays.asList( firstSupplier, secondSupplier, thirdSupplier ); // 构建链式Observable并订阅 chainSuppliers(supplierList).subscribe(); // 测试触发流程 firstSubject.onNext(""); // 此时才会打印side effect of second one secondSubject.onNext(""); // 此时才会打印side effect of third one
逻辑说明
reduce的作用是把独立的供应商合并成一个嵌套的供应商:每一步合并都会生成新的供应商,它的get()会先执行前一个供应商的get()拿到Observable,再通过flatMap监听这个Observable的onNext事件,只有事件触发时,才会调用下一个供应商的get()。- 完全满足延迟需求:所有供应商的
get()都不会提前执行,只有当前面的Observable发射数据后,才会触发下一个。
方法二:纯RxJava操作符实现
如果你更倾向于用RxJava原生操作符,也可以这样写:
public static <T> Observable<T> chainSuppliers(Iterable<Supplier<Observable<T>>> suppliers) { return Observable.fromIterable(suppliers) .reduce(Observable.empty(), (previousObservable, currentSupplier) -> // 处理第一个Observable的特殊情况 previousObservable.isEmpty() ? currentSupplier.get() : previousObservable.flatMap(ignored -> currentSupplier.get()) ) .flatMap(observable -> observable); }
这个版本用Observable.fromIterable把供应商集合转成Observable流,再用reduce逐步构建链式结构,逻辑和第一种方法完全一致,更贴合RxJava的使用习惯。
内容的提问来源于stack exchange,提问作者Piotr Aleksander Chmielowski
相关产品推荐
相关产品推荐

