如何在RxJava中等待多个Observable执行完成?
解决方案:用RxJava原生流替代PublishSubject实现结果逐个传递+自动完成
核心思路
抛弃手动管理PublishSubject和任务计数的方式,直接通过RxJava的流操作将所有任务整合成一个Observable序列:
- 将
operations集合转为Observable流,自动绑定每个任务的ID - 每个任务在独立线程执行,结果携带ID后发送到主线程
- 整个流的
onComplete会在所有任务执行完毕后自动触发,无需手动计数
代码重构
1. 改造任务执行方法(返回Observable而非直接订阅)
private Observable<Pair<Integer, Long>> performOperation(Operation operation, int id, int operationsAmount) { return Observable.defer(() -> operation.executeAndReturnUptime(operationsAmount)) .subscribeOn(Schedulers.newThread()) // 每个任务在新线程执行 .map(upTime -> new Pair<>(id, upTime)); // 给结果绑定ID }
2. 重构主方法生成完整Observable流
public Observable<Pair<Integer, Long>> createAndPerformOperations(int mDataStructureSize, int operationsAmount) { return Single.fromCallable(() -> // 生成任务集合(原逻辑不变) new OperationsFactory().getOperations(new DataStructureFactory().getMaps(mDataStructureSize)) ) .subscribeOn(Schedulers.newThread()) // 任务集合生成也在后台线程 .flatMapObservable(operations -> // 将任务集合转为Observable,同时绑定索引作为ID Observable.fromIterable(operations) .zipWith(Observable.range(0, operations.size()), (operation, id) -> performOperation(operation, id, operationsAmount) .observeOn(AndroidSchedulers.mainThread()) // 结果回到主线程 ) .flatMap(observable -> observable) // 展开所有任务的Observable流 ); }
3. 在Fragment中订阅流
// 替换原来对uptimeStream的订阅,直接订阅这个Observable disposables.add(viewModel.createAndPerformOperations(mDataStructureSize, operationsAmount) .subscribe( pair -> { // 逐个处理每个任务的结果(id和upTime) int taskId = pair.first; long upTime = pair.second; // 你的业务逻辑 }, throwable -> { // 处理任务执行中的错误 }, () -> { // 所有任务完成且结果都已消费,这里执行原onComplete逻辑 } ) );
关键优势
- 无需手动计数:整个流的
onComplete由RxJava自动管理,所有任务执行完毕后自动触发 - 避免PublishSubject的副作用:用原生Observable流管理结果传递,更符合RxJava的响应式设计
- 灵活适配任务数量变化:不管任务数量多少,流都会自动处理,无需修改代码
- 解决ID绑定问题:通过
zipWith(Observable.range(...))自动给每个任务分配ID,无需手动循环索引
内容的提问来源于stack exchange,提问作者Andrei Aleksandrov
相关产品推荐
相关产品推荐

