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
相关产品推荐
相关产品推荐

