【问题标题】:RxJava - how to wait to result of async tasks wihiting doOnNextRxJava - 如何使用 doOnNext 等待异步任务的结果
【发布时间】:2021-09-16 10:41:07
【问题描述】:

如何实现 doOnNext 等待多个异步任务的结果?

例如-

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(user -> {
                ArrayList<String> imagesUrls = user.getUrls();
                for (String url : imagesUrls) {
                    storage.getReference().child("images").child(url).getBytes(ParametersConventions.FIREBASE_DOWNLOAD_IMAGE_MAX_SIZE).
                    addOnSuccessListener(bytes -> {
                       doSomething(bytes);    
                    });
                }
            })
            .doOnNext(user -> {
                doSomething();
            })
            .doOnComplete(...);
}

我希望在所有下载图像的异步调用完成后调用 doSomething 的 doOnNext。

【问题讨论】:

    标签: android rx-java


    【解决方案1】:

    将该 API 调用转换为响应式类型并将其合并到主流程中:

    int max = ParametersConventions.FIREBASE_DOWNLOAD_IMAGE_MAX_SIZE;
    
    public Completable downloadAsync(URL url) {
        return Completable.create(inner -> {
                               storage.getReference()
                               .child("images")
                               .child(url)
                               .getBytes(max)
                               .addOnSuccessListener(bytes -> {
                                    doSomething(bytes);
                                    inner.onComplete();
                               });
                         });
    }
    

    一起:

    
    Flowable.create(emitter-> {
           emitter.onNext(user);
        }, BackpressureStrategy.BUFFER)
                .observeOn(Schedulers.io())
                .concatMapSingle(user -> 
                     Flowable.fromIterable(user.getUrls())
                         .concatMapCompletable(url -> downloadAsync(url))
                         .andThen(Single.just(user))
                 )
                 .doOnNext(user -> {
                     doSomething();
                 })
                 .doOnComplete(...);
    

    【讨论】:

    • 非常感谢,这是工作。但也许您知道是否可以并行运行所有 Completable?似乎它一个接一个地运行 Completable,所以它比较慢。
    • flatMapCompletable.
    【解决方案2】:

    doOnNext 运算符会在每次流中有新项目时触发,因此它不是您的最佳选择。根据您的需要尝试使用 map/flatMap/concatMap 运算符。如果您需要打几个电话,然后对数据做一些事情,您可以查看我已经回答过的类似问题链接:Chaining API Requests with Retrofit + Rx 您可以在其中找到一种方法来进行顺序网络调用,然后对数据列表执行任何操作:D

    【讨论】:

    • 我遇到的问题是我用于异步任务的数据不是流经流的对象,所以我不能像我理解的那样使用地图的运算符。我需要找到一种方法来阻止异步任务或以某种方式使用回调。
    • 将 asyncTask 与 Rxjava 混合并不是最好的做法。我建议您只使用 RxJava,这将提高代码的可读性和可测试性。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-10-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-03-18
    相关资源
    最近更新 更多