【问题标题】:Do something after multiple observables have completed在多个 observables 完成后做某事
【发布时间】:2016-10-21 02:29:49
【问题描述】:

我在 Android 上使用 RXJava 并尝试将多个 API 调用链接在一起,并在两个 API 调用完成后执行一些操作。我的 API 调用看起来都与提供的代码示例相似。基本上是做API调用,在onNext中将每条记录写入DB,等所有记录都写入后,更新一些缓存。我想异步触发这两个调用,然后在两者都点击 onCompleted 之后,然后做其他事情。 RX 中执行此操作的正确方法是什么?我认为我不需要 zip,因为我不需要将不同的流捆绑在一起。我在想也许可以合并,但我的两个 API 调用返回了不同类型的 Observable。请告诉我。谢谢。

    getUsers()
            .subscribeOn(Schedulers.io())
            .observeOn(Schedulers.io())
            .flatMap(Observable::from)
            .subscribe(new Subscriber<User>() {
                @Override
                public void onCompleted() {
                   updateUserCache();
                }

                @Override
                public void onError(Throwable e) {
                    Log.e(TAG, "Error loading users", e);
                }

                @Override
                public void onNext(User user) {
                    insertUserToDB(user);
                }
            });

    getLocations()
            .subscribeOn(Schedulers.io())
            .observeOn(Schedulers.io())
            .flatMap(Observable::from)
            .subscribe(new Subscriber<Location>() {
                @Override
                public void onCompleted() {
                   updateLocationCache();
                }

                @Override
                public void onError(Throwable e) {
                    Log.e(TAG, "Error loading Locations", e);
                }

                @Override
                public void onNext(Location location) {
                    insertLocationToDB(location);
                }
            });  

【问题讨论】:

    标签: java android rx-java reactive-programming rx-android


    【解决方案1】:

    你的想法是正确的。您应该使用zip 运算符。

    您的每个函数都应该进行调用、写入数据库并执行您需要的所有操作。 Theat zip 的输出功能不同:当它被调用时,您可以确定所有 Observable 都已成功完成 -> 只需完成您的反应流。

    创建Observable的列表:

    List<Observable<?>> observableList = new ArrayList<>();
    observableList.add(
            getUsers()
                .subscribeOn(Schedulers.io())
                .observeOn(Schedulers.io())
                .flatMap(Observable::from)
                .insertUserToDB(user)
                .toList());
    
    observableList.add(
            getLocations()
                .subscribeOn(Schedulers.io())
                .observeOn(Schedulers.io())
                .flatMap(Observable::from)
                .insertLocationToDB(location)
                .toList());
    

    然后zip alll Observable的:

    Observable.zip(observableList, new FuncN<Object, Observable<?>>() {
        @Override
        public Observable<?> call(Object... args) {
            return Observable.empty();
        }
    }).subscribe(new Subscriber<Object>() {
        @Override
        public void onCompleted() {
            updateUserCache();
            updateLocationCache();
        }
    
        @Override
        public void onError(Throwable e) {
    
        }
    
        @Override
        public void onNext(Object o) {
    
        }
    });
    

    这是伪代码,但我希望你能理解这个想法。

    【讨论】:

    • 感谢您的回复。几个问题。首先对 insertUserToDB 的调用,不需要在 subscribe 方法调用中发生,这意味着我不能为 toList() 返回一个 observable?
    • 我最终使用 doOnNext 插入数据库。
    【解决方案2】:

    如果有人需要,这是我根据 R. Zagórski 建议使用的代码:

        List<Observable<?>> observableList = new ArrayList<>();
        observableList.add(
                getUsers()
                    .subscribeOn(Schedulers.io())
                    .observeOn(Schedulers.io())
                    .flatMap(Observable::from)
                    .doOnNext(user->insertUser(user))
                    .toList()
        );
        observableList.add(
                getLocations()
                        .subscribeOn(Schedulers.io())
                        .observeOn(Schedulers.io())
                        .flatMap(Observable::from)
                        .doOnNext(location->insertLocation(location))
                        .toList()
        );
    
        Observable.zip(observableList, new FuncN<Object>() {
            @Override
            public Observable<?> call(Object...args) {
                return Observable.empty();
            }).subscribe(new Subscriber<Object>() {
                @Override
                public void onCompleted() {
                    updateUserCache();
                    updateLocationCache();
                }
    
                @Override
                public void onError(Throwable e) {
    
                }
    
                @Override
                public void onNext(Object o) {
    
                }
            });
    

    【讨论】:

      【解决方案3】:

      .zip() 是正确的方法

      你可能想让 Retrofit 返回 Single 而不是 Observable

      【讨论】:

        猜你喜欢
        • 2020-10-27
        • 1970-01-01
        • 2018-07-24
        • 1970-01-01
        • 2016-09-06
        • 1970-01-01
        • 2014-07-24
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多