【问题标题】:RxJava pattern for returning cold results, doing more work, then returning hot results用于返回冷结果、做更多工作、然后返回热结果的 RxJava 模式
【发布时间】:2016-11-07 12:20:20
【问题描述】:

我正在学习 RxJava,所以请温柔一点。我看过教程,阅读过,搜索过,但是,我在转换我的AsyncTaskLoader 时仍然遇到一些问题。出于某种原因,我找不到一种运算符模式来完成我的任务(尽管我认为这是一种常见的模式)。我要做的是:返回一个我的片段可以订阅的 Observable。 observable 应该在订阅时执行以下操作:

1) 通过执行 2 次查询、运行一些逻辑并返回结果从本地数据库中获取数据;
2) 从 API 获取数据;
3) 将新的 API 数据与数据库同步;
4) 重复第一步并返回结果;

到目前为止,我已经转换了我的数据库调用和我的 API 调用以返回 observables。我试图了解如何发出冷结果并继续使用链条。我可能可以分别保留这两个操作,并使用同一个订阅者来订阅两者?但是我不确定如果我的新 loader-replacement 类返回一个 observable 它将如何工作......而且我真的不需要处理第二个 observable 的结果 - 我只需要第一个在第二个完成。

到目前为止,我有以下内容:

 public Observable<StuffFetchResult> getColdStuff() {
    return Observable.zip(mDataSource.listStuff(), mDataSource.listOtherStuff(),
            (stuff, moreStuff) -> {
                List<Stuff> mergedList = new ArrayList<>();
                // do some merging stuff
                return new StuffFetchResult(mergedList);
            }).subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread());
 } 

假设我也有 getHotStuff(),它将执行 API 调用并与数据库同步,如果这是正确的方法,并返回相同的 Observable。但是,我被困在下一步 - 一旦hotStuff 完成后,如何重新启动第一个可观察的重播,而不添加其他订阅者?

编辑:

我已经取得了一些进展,我认为我现在需要的只是加入这一切。我有两种方法:

1) getColdStuff() 和上面描述的差不多
2)getHotStuff() 会调用 API,与数据库同步,并返回一个 Observable。想法是在getHotStuff() 完成后再次调用getColdStuff() 以刷新UI,因此可以忽略从getHotStuff() 返回的实际结果。它所要做的就是在完成后触发getColdStuff()

我已经尝试了答案中的建议并创建了以下内容:

BehaviorRelay<Observable<StuffFetchResult>> callSequence = BehaviorRelay.create();
Observable<StuffFetchResult> valueSequence = Observable.switchOnNextDelayError(callSequence.toSerialized());
valueSequence.subscribe(new Subscriber<StuffFetchResult>() {
    @Override
    public void onCompleted() {}

    @Override
    public void onError(Throwable e) {}

    @Override
    public void onNext(StuffFetchResult result) {
        // UI stuff
    }
});
callSequence.call(loader.getColdStuff());

我可以在这里订阅valueSequence 并使用callSequence.call(loader.getColdStuff());,它将运行第一个方法并在我的订阅的onNext() 中产生结果,我可以将其用于我的UI。但是,我不确定如何并行运行 getHotStuff() 并在它返回时对其执行不同的操作。另外getHotStuff() 返回不同类型的Observable,所以我不能真正使用相同的callSequence

编辑 2

使用两个订阅者,我可以实现我认为的所需行为。不确定这是否是正确的方法。

 loader.getHotStuff()
    .subscribeOn(Schedulers.io())
    .subscribe( new Subscriber<Object>() {
        @Override
        public void onCompleted() {}

        @Override
        public void onError(Throwable e) {}

        @Override
        public void onNext(Object stuffWeDontCareAbout) {
            callSequence.call(loader.getColdStuff());
        }
    });

【问题讨论】:

    标签: android design-patterns rx-java


    【解决方案1】:

    如果我正确理解你的场景,你可能想要这样的东西 -

    BehaviorSubject<Observable<T> callSequence = BehaviorSubject.create();
    Observable<T> valueSequence = Observable.swithOnNextDelayError(callSequence.toSerialized());
    

    您的订阅者将监听 valueSequence,并且当您需要“重新启动”时,您将调用它 -

    callSequence.onNext(call.cache()); // *call* is Observable<T>
    

    (我将 .subscribeOn/.observeOn 配置留给你)

    【讨论】:

    • 感谢您的回答。您能否详细说明一下,因为我不确定如何使用您的建议?此外,我注意到我的方法 getHotStuffgetColdStuff 返回不同类型的 observables - 因为最初 getHotStuff 将从 API 获取,持久化,然后调用 getColdStuff 从存储中检索。所以大概他们不能都在callSequence?我正在寻找的行为是:getCold -&gt; resultsToUI -&gt; getHot -&gt; persist -&gt; getCold - resultsToUI
    • 对于副作用,例如持久化到存储,.doOnNext() 是合适的地方。如果第一个为空,则使用备用数据源,.switchIfEmpty(alternate) 是典型的解决方案。所以整个链条就像 storage.switchIfEmpty(apiCall.doOnNext(persist()));不同类型的 observable 几乎不是问题,您可以使用 .map()
    • 再次感谢。我仍在尝试将其放在一起。这并不是说我需要使用替代数据源,实际上恰恰相反-我需要同时使用两者,首先显示来自陈旧源的数据,然后在我有新数据时进行更新。但我不确定如何将它们全部连接到一个流程中。我会用最新的代码更新我的问题,以防你还有一分钟的时间看看。谢谢
    • 我想你想尽快从 API 开始更新,即并行。所以,我首先将 UI 订阅到“valueSequence”,然后调用“callSequence.onNext(hotStuff.doOnNext(/*persist*/).flatMap(/*to coldStuff*/).cache());”。这样,您将在 API 调用前后更新 UI
    • 我认为这几乎可行。我让热的东西做它的事情,我让冷的东西在热的东西之后运行,几乎是神奇的。我现在唯一的问题是初始的coldStuff 没有运行,理想情况下我希望它与hotStuff 并行运行。但是,如果我尝试向callSequence 添加两个调用,则只有最后一个调用通过。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-05-15
    • 1970-01-01
    • 2016-02-05
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多