【问题标题】:Android: infinite scroll with rx-java using repeatWhen, takeUntil and filter with retrofitAndroid:使用repeatWhen,takeUntil和过滤器与rx-java无限滚动
【发布时间】:2017-03-01 16:37:59
【问题描述】:

我正在使用带有 RxJava 的 Retrofit 2.2。 分页是这样工作的:我得到第一批数据,我必须请求第二批具有相同参数的数据,除了一个是 lastUpdated 日期,然后如果我得到空或同一批数据,则意味着有没有更多的项目。我找到了这篇很棒的文章https://medium.com/@v.danylo/server-polling-and-retrying-failed-operations-with-retrofit-and-rxjava-8bcc7e641a5a#.40aeibaja,关于如何做到这一点。所以我的代码是:

private Observable<Integer> syncDataPoints(final String baseUrl, final String apiKey,
        final long surveyGroupId) {
    final List<ApiDataPoint> lastBatch = new ArrayList<>();
    Timber.d("start syncDataPoints");
    return loadAndSave(baseUrl, apiKey, surveyGroupId, lastBatch)
            .repeatWhen(new Func1<Observable<? extends Void>, Observable<?>>() {
                @Override
                public Observable<?> call(final Observable<? extends Void> observable) {
                    Timber.d("Calling repeatWhen");
                    return observable.delay(5, TimeUnit.SECONDS);
                }
            })
            .takeUntil(new Func1<List<ApiDataPoint>, Boolean>() {
                @Override
                public Boolean call(List<ApiDataPoint> apiDataPoints) {
                    boolean done = apiDataPoints.isEmpty();
                    if (done) {
                        Timber.d("takeUntil : finished");
                    } else {
                        Timber.d("takeUntil : will query again");
                    }
                    return done;
                }
            })
            .filter(new Func1<List<ApiDataPoint>, Boolean>() {
                @Override
                public Boolean call(List<ApiDataPoint> apiDataPoints) {
                    boolean unfiltered = apiDataPoints.isEmpty();
                    if (unfiltered) {
                        Timber.d("filtered");
                    } else {
                        Timber.d("not filtered");
                    }
                    return unfiltered;
                }
            }).map(new Func1<List<ApiDataPoint>, Integer>() {
                @Override
                public Integer call(List<ApiDataPoint> apiDataPoints) {
                    Timber.d("Finished polling server");
                    return 0;
                }
            });
}

private Observable<List<ApiDataPoint>> loadAndSave(final String baseUrl, final String apiKey,
        final long surveyGroupId, final List<ApiDataPoint> lastBatch) {
    return loadNewDataPoints(baseUrl, apiKey, surveyGroupId)
            .concatMap(new Func1<ApiLocaleResult, Observable<List<ApiDataPoint>>>() {
                @Override
                public Observable<List<ApiDataPoint>> call(ApiLocaleResult apiLocaleResult) {
                    return saveToDataBase(apiLocaleResult, lastBatch);
                }
            });
}


private Observable<ApiLocaleResult> loadNewDataPoints(final String baseUrl, final String apiKey,
        final long surveyGroupId) {
    Timber.d("loadNewDataPoints");

    return Observable.just(true).concatMap(new Func1<Object, Observable<ApiLocaleResult>>() {
        @Override
        public Observable<ApiLocaleResult> call(Object o) {
            Timber.d("loadNewDataPoints call");
            return restApi
                    .loadNewDataPoints(baseUrl, apiKey, surveyGroupId,
                            getSyncedTime(surveyGroupId));
        }
    });
}

如您所见,有趣的方法是loadNewDataPoints,我希望在没有更多数据点之前调用它。如您所见,Observable.just(true).concatMap 是一个 hack,因为如果我删除此 concat 映射,restApi.loadNewDataPoints(....) 不会被调用,尽管在日志中我可以看到 api 确实被调用但使用相同的旧参数,当然它返回结果与第一次相同,因此同步停止, saveToDataBase 确实可以正常调用。使用我的 hack 它可以工作,但我想了解为什么它不能以另一种方式工作,以及是否有更好的方法来做到这一点。非常感谢!

【问题讨论】:

    标签: android rx-java retrofit2 rx-android


    【解决方案1】:

    所以,我编写了这种 API(称为 Keyset Pagination)并针对它们实现了 Rx 客户端。

    这是 BehaviorSubjects 有用的情况之一:

    S initialState = null;
    BehaviorProcessor<T> subject = BehaviorProcessor.createDefault(initialState);
    return subject
      .flatMap(state -> getNextElements(state).singleOrError().toFlowable(), Pair::of, 1)
      .serialize()
      .flatMap(stateValuePair -> {
          S state = stateValuePair.getLeft();
          R retrievedValue = stateValuePair.getRight();
          if(isEmpty(retrievedValue)) {
             subject.onComplete();
             return Flowable.empty();
          } else {
             subject.onNext(getNextState(state, retrievedValue));
             return Flowable.just(retrievedValue);
          }
        }
       .doOnUnsubscribe(subject::onCompleted)
       .map(value -> ...)
    

    这里

    • getNextElement 基于状态执行网络调用并返回具有单个值的反应流
    • isEmpty判断返回值是否为空表示元素结束
    • getNextState 将传入的状态与检索到的值相结合,以确定 getNextElement 的下一个状态。

    如果发生错误(它将被传播)并且如果您在结束前取消订阅(查询将被终止),它将正常工作。

    当然,在您的具体情况下,这些不需要是单独的方法或复杂类型。

    【讨论】:

    • 感谢您对我搜索过的键集分页的解释,它似乎比偏移量use-the-index-luke.com/no-offset 更好。还要感谢我将尝试这个的代码,我还没有使用 Flowables。 “这些不需要是单独的方法或复杂类型”是什么意思?
    • 我的意思是代码段更像伪代码 - 你可以有 f.e if(retrievedValue.getList().isEmpty())。此外,你也可以使用 Observables,但那样你就会失去背压。
    猜你喜欢
    • 2015-01-29
    • 1970-01-01
    • 2013-12-01
    • 2021-06-14
    • 1970-01-01
    • 2013-10-31
    • 1970-01-01
    • 2020-11-28
    • 2021-07-05
    相关资源
    最近更新 更多