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

RxJava如何让doOnNext等待多个异步任务全部执行完成后再执行后续逻辑

RxJava 等待 doOnNext 内部多个异步任务完成的解决方案

原代码的问题根源在于:doOnNext 仅会执行你传入的 lambda 代码块,不会等待内部异步回调完成。你在 doOnNext 中发起的 Firebase 下载是异步回调式任务,遍历完所有 URL 后 doOnNext 就会立刻放行,下游逻辑自然不会等待下载完成。

改造思路如下:

  • 第一步:将 Firebase 回调式下载封装为 RxJava Single<byte[]> 类型,让 RxJava 可以感知异步任务的生命周期
  • 第二步:用 flatMap 替换第一个 doOnNext,批量执行所有下载任务,等待全部完成后再向下游发射 User 对象

代码实现

首先封装下载方法为RxJava类型:

// 封装Firebase下载为Single,让RxJava感知异步任务状态
private Single<byte[]> downloadImage(String url) {
    return Single.create(emitter -> {
        storage.getReference().child("images").child(url)
                .getBytes(ParametersConventions.FIREBASE_DOWNLOAD_IMAGE_MAX_SIZE)
                .addOnSuccessListener(bytes -> {
                    doSomething(bytes);
                    if (!emitter.isDisposed()) {
                        emitter.onSuccess(bytes);
                    }
                })
                .addOnFailureListener(e -> {
                    if (!emitter.isDisposed()) {
                        emitter.tryOnError(e);
                    }
                });
    });
}

改造后的核心流逻辑:

public void getImages(User user) {
    Flowable.create(new FlowableOnSubscribe<User>() {
                @Override
                public void subscribe(@io.reactivex.rxjava3.annotations.NonNull FlowableEmitter<User> emitter) throws Throwable {
                    emitter.onNext(user);
                }
            }, BackpressureStrategy.BUFFER)
            .observeOn(Schedulers.io())
            // 替换第一个doOnNext,等待所有下载任务完成再向下游发射数据
            .flatMap(u -> {
                ArrayList<String> imagesUrls = u.getUrls();
                // 把所有下载任务转成Completable列表,只关注任务是否完成
                List<Completable> downloadTasks = imagesUrls.stream()
                        .map(this::downloadImage)
                        .map(Single::ignoreElement)
                        .toList();
                // 合并所有任务,全部完成后发射原User对象到下游
                return Completable.merge(downloadTasks)
                        .andThen(Flowable.just(u));
                // 若需要忽略单个下载失败,不中断整个流,可替换为以下代码:
                // return Completable.merge(downloadTasks.stream().map(task -> task.onErrorComplete()).toList())
                //         .andThen(Flowable.just(u));
            })
            // 执行到此处时,所有图片下载任务已全部完成
            .doOnNext(u -> {
                doSomething();
            })
            .doOnComplete(...)
            .subscribe(); // 注意必须添加订阅,流才会执行
}

可选优化说明

  • 如果需要收集所有下载的字节结果,可以去掉 ignoreElement() 调用,用 toList() 收集所有下载结果后和原User一起发射到下游
  • 可根据业务需求给单个下载任务添加重试、降级逻辑,不会影响整体流的执行

内容的提问来源于stack exchange,提问作者nirkov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 22:24:00