You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在RxJava中等待多个Observable执行完成?

解决方案:用RxJava原生流替代PublishSubject实现结果逐个传递+自动完成

核心思路

抛弃手动管理PublishSubject和任务计数的方式,直接通过RxJava的流操作将所有任务整合成一个Observable序列:

  1. 将operations集合转为Observable流,自动绑定每个任务的ID
  2. 每个任务在独立线程执行,结果携带ID后发送到主线程
  3. 整个流的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 13:50:16