【问题标题】:Executing rx.Obseravables secuentially顺序执行 rx.Observables
【发布时间】:2017-04-20 15:04:11
【问题描述】:

我正在使用 Fernando Ceja 的简洁架构开发一个 Android 应用。我的交互者或用例之一负责获取用户的提要数据。为了获取数据,首先我必须从数据库表中检索用户的团队,然后我必须从服务器端获取 Feed 列表。

这就是我从数据库层获取团队的方式:

    mTeamCache.getAllTeams().subscribe(new DefaultSubscriber<List<SimpleTeam>>() {
        @Override
        public void onNext(List<SimpleTeam> simpleTeams) {
            super.onNext(simpleTeams);
            mTeams = simpleTeams;
        }
    });

TeamCache 基本上只是另一个交互器,负责获取我在数据库中拥有的所有团队。

以下是我从服务器端获取 Feed 数据的方式:

    mFeedRepository.getFeed(0, 50).subscribe(new ServerSubscriber<List<ApiFeedResponse>>() {
        @Override
        protected void onServerSideError(Throwable errorResponse) {
            callback.onFeedFetchFailed(...);
        }

        @Override
        protected void onSuccess(List<ApiFeedResponse> responseBody) {
            //Do stuff with mTeams
            callback.onFeedFetched(...);
        }
    });

我的 GetFeedInteractor 类有一个名为 execute 的方法,我通过回调传递我稍后在 UI 中使用它来处理响应。所有这一切的问题是,目前我正在链接这样的响应:

@Override
public void execute(final Callback callback, String userSipId) {
    mTeamCache.getAllTeams().subscribe(new DefaultSubscriber<List<SimpleTeam>>() {
        @Override
        public void onNext(List<SimpleTeam> simpleTeams) {
            super.onNext(simpleTeams);
            mTeams = simpleTeams;
            getFeedFromRepository(callback);
        }
    });
}

public void getFeedFromRepository(final Callback callback) {
    mFeedRepository.getFeedRx(0, 50).subscribe(new ServerSubscriber<List<ApiFeedResponse>>() {
        @Override
        protected void onServerSideError(Throwable errorResponse) {
            callback.onFeedFetchFailed("failed");
        }

        @Override
        protected void onSuccess(List<ApiFeedResponse> responseBody) {
            //Do stuff with mTeams
            List<BaseFeedItem> responseList = new ArrayList();

            for (ApiFeedResponse apiFeedResponse : responseBody) {
                responseList.add(FeedDataMapper.transform(apiFeedResponse));
            }

            callback.onFeedFetched(responseList);
        }
    });
}

如您所见,一旦我从缓存交互器获取团队集合,我就会调用从同一个订阅者获取提要的方法。我不喜欢这个。我希望能够做一些更好的事情,比如使用 Observable.concat(getTeamsFromCache(), getFeedFromRepository());在订阅者中将调用链接到另一个 rx.Observable 并不是一件好事。我想我的问题是,如何链接两个使用不同订阅者的 rx.Observables?

更新:

ServerSubscriber 是我实现订阅改造服务的订阅者。它只是检查错误代码和一些东西。这里是:

https://gist.github.com/4gus71n/65dc94de4ca01fb221a079b68c0570b5

默认订阅者是一个空的默认订阅者。这里是:

https://gist.github.com/4gus71n/df501928fc5d24c2c6ed7740a6520330

TeamCache#getAllTeams() 返回 rx.Observable> FeedRepository#getFeed(int page, int offset) 返回 rx.Observable>

更新 2:

这是获取用户提要的交互器现在的样子:

@Override
    public void execute(final Callback callback, int offset, int pageSize) {
        User user = mGetLoggedUser.get();
        String userSipid = mUserSipid.get();

        mFeedRepository.getFeed(offset, pageSize) //Get items from the server-side
                .onErrorResumeNext(mFeedCache.getFeed(userSipid)) //If something goes wrong take it from cache
                .mergeWith(mPendingPostCache.getAllPendingPostsAsFeedItems(user)) //Merge the response with the pending posts
                .subscribe(new DefaultSubscriber<List<BaseFeedItem>>() {
                    @Override
                    public void onNext(List<BaseFeedItem> baseFeedItems) {
                        callback.onFeedFetched(baseFeedItems);
                    }

                    @Override
                    public void onError(Throwable e) {
                        if (e instanceof ServerSideException) {
                            //Handle the http error
                        } else if (e instanceof DBException) {
                            //Handle the database cache error
                        } else {
                            //Handle generic error
                        }
                    }
                });
    }

【问题讨论】:

  • 什么是 getFeedRx(0, 50) ?你在这里 Observable 吗? ServerSubscriber 是什么?
  • 使用该信息更新帖子
  • 那里,我添加了更多信息

标签: android rx-java clean-architecture


【解决方案1】:

我认为你错过了 RxJava 和响应式方法的要点,你不应该有不同的订阅者与 OO 层次结构和回调。
您应该构造分离的Observables,它应该发出它所处理的特定数据,没有Subscriber,然后您可以根据需要链接Observable,最后,您拥有对最终结果做出反应的订阅者预计来自链接的 Observable 流。

类似这样的东西(使用 lambdas 来编写更精简的代码):

TeamCache mTeamCache = new TeamCache();
FeedRepository mFeedRepository = new FeedRepository();

Observable.zip(teamsObservable, feedObservable, Pair::new)
    .subscribe(resultPair -> {
            //Do stuff with mTeams
            List<BaseFeedItem> responseList = new ArrayList(); 
            for (ApiFeedResponse apiFeedResponse : resultPair.second) {
                responseList.add(FeedDataMapper.transform(apiFeedResponse));
            }
        }, throwable -> {
             //handle errors
           }
     );

我使用 zip 而不是 concat,因为您似乎在这里有 2 个独立的调用,您希望等待两者完成(将它们“压缩”在一起)然后采取行动,但当然,当您已分离Observables 流,您可以根据需要将它们以不同的方式链接在一起。

至于您的 ServerSubscriber 以及所有响应验证逻辑,它也应该是 rxify,因此您可以将其组合到您的服务器 Observable 流中。

类似的东西(为了简化而发出的一些逻辑,因为我不熟悉它......)

Observable<List<SimpleTeam>> teamsObservable = mTeamCache.getAllTeams();
Observable<List<ApiFeedResponse>> feedObservable = mFeedRepository.getFeed(0, 50)
            .flatMap(apiFeedsResponse -> {
                if (apiFeedsResponse.code() != 200) {
                    if (apiFeedsResponse.code() == 304) {
                        List<ApiFeedResponse> body = apiFeedsResponse.body();
                        return Observable.just(body);
                        //onNotModified(o.body());
                    } else {
                        return Observable.error(new ServerSideErrorException(apiFeedsResponse));
                    }
                } else {
                    //onServerSideResponse(o.body());
                    return Observable.just(apiFeedsResponse.body());
                }
            });

【讨论】:

  • 不错的方法。我喜欢使用 flatMap 处理服务器端错误。只是为了阅读更多意见,我会在结束这个问题之前等待一段时间。
  • 顺便说一句,我正在尝试执行与 flatMap 相同的操作来映射 http 错误。但我想通过使用自定义 rx.Observables 和更有用的 API 来改进它。只有 onNext、onComplete 和 onError 不足。您知道如何扩展 rx.Observable 或 rx.Subscriber 以向这些对象的“生命周期”添加更多方法吗?我正在搜索,但到目前为止我发现的唯一有用的是这个github.com/ReactiveX/RxJava/issues/1034
  • 我不敢苟同,试图扩展 Observable/Subscriber 会带来麻烦,因为它的强大之处在于它的通用 API 可以适合任何东西,从而为您提供强大的编写能力,链接等此外添加“生命周期”将回到回调世界。我认为 a) 反应式让您的想法与 OO 不同 b) 您仍然具有扩展的能力,您可以将复杂的状态封装在将要发出的某个对象中,然后将所有状态放在 onNext。
  • 我明白你的意思。但检查我的更新。我需要处理 ServerSideExceptions,如果使用连接到 Retrofit Service 调用的转换运算符出现任何问题,我将抛出该异常,并且我需要处理 DBExceptions,以防在调用数据库缓存时出现问题。我不喜欢那些instanceof。现在我想我要创建一个类似 MyAppSubscriber 的东西,在那个类中执行那些 instanceof,并调用几个名为 onServerSideError() 和 onDbError() 的方法。您认为这没问题还是有更好的方法?
  • 这是一个正确的问题,我建议你打开关于这个问题的附加问题,我可以更详细地回答一个答案
猜你喜欢
  • 2012-11-23
  • 2018-08-12
  • 2013-11-12
  • 2012-04-10
  • 2013-04-11
  • 1970-01-01
  • 2017-02-10
  • 2011-04-12
  • 1970-01-01
相关资源
最近更新 更多