RxJava2如何并行执行多个Completable并监听全部完成或错误
嘿,这个问题我太熟了!你现在用的Completable.concat()是串行执行任务,要改成并行的话,RxJava里有现成的操作符完美适配你的需求——不用纠结Observable.zip,因为Completable家族本身就有专门的合并操作~
核心方案:用Completable.merge()实现并行合并
Completable.merge()(或者Completable.mergeArray(),适合固定数量的任务)就是专门用来并行执行多个Completable的操作符,它的行为完全符合你的要求:
- 所有传入的Completable会同时并行执行
- 当全部任务都成功完成时,才会触发
onComplete回调 - 只要任意一个任务抛出异常,会立刻触发
onError回调,并且自动取消其他正在运行的任务(如果需要等待所有任务完成再抛出错误,可以用mergeDelayError,不过你的需求显然更适合前者)
直接把你的代码改成这样就行:
Completable a = computationTaskA(); Completable b = computationTaskB(); Completable c = computationTaskC(); // 用merge替代concat,实现并行执行 Completable all = Completable.merge(Arrays.asList(a, b, c)) .subscribe( () -> { // 所有任务都成功完成,处理完成逻辑 }, error -> { // 任意任务失败,处理错误逻辑 } );
如果任务数量固定,也可以用更简洁的mergeArray:
Completable all = Completable.mergeArray(a, b, c) .subscribe(...);
关键细节:给阻塞任务指定调度器
这里要注意一个坑:如果你的computationTaskA()这类方法是阻塞式的(比如网络请求、 heavy计算),一定要给每个Completable加上subscribeOn(),指定合适的调度器,否则所有任务可能还是会在同一个线程串行执行,达不到并行效果。
比如计算任务用Schedulers.computation(),网络请求用Schedulers.io():
Completable a = computationTaskA() .subscribeOn(Schedulers.computation()); // 让计算任务在计算线程池执行 Completable b = networkTaskB() .subscribeOn(Schedulers.io()); // 网络请求用IO线程池 Completable c = computationTaskC() .subscribeOn(Schedulers.computation()); Completable all = Completable.merge(Arrays.asList(a, b, c)) .observeOn(AndroidSchedulers.mainThread()) // 如果是Android平台,切换回主线程处理回调 .subscribe(...);
为什么不用Observable.zip?
你提到Observable.zip()不适用,这点是对的——因为zip需要每个Observable发射至少一个数据,但Completable根本不发射数据,只通知完成/错误。而Completable.merge()完全针对无结果的任务设计,完美匹配你的场景。
内容的提问来源于stack exchange,提问作者ddinchev
相关产品推荐
相关产品推荐

