【问题标题】:Paginate Observable results without recursion - RxJava无需递归的分页 Observable 结果 - RxJava
【发布时间】:2016-09-16 12:29:06
【问题描述】:

我有一个非常标准的 API 分页问题,​​您可以通过一些简单的递归来处理它。这是一个虚构的例子:

public Observable<List<Result>> scan() {
    return scanPage(Optional.empty(), ImmutableList.of());
}

private Observable<?> scanPage(Optional<KEY> startKey, List<Result> results) {
    return this.scanner.scan(startKey, LIMIT)
            .flatMap(page -> {
                if (!page.getLastKey().isPresent()) {
                    return Observable.just(results);
                }
                return scanPage(page.getLastKey(), ImmutableList.<Result>builder()
                        .addAll(results)
                        .addAll(page.getResults())
                        .build()
                );
            });
}

但这显然会创建一个庞大的调用堆栈。我怎样才能强制执行此操作但保持 Observable 流?

这是一个命令式阻塞示例:

public List<Result> scan() {
    Optional<String> startKey = Optional.empty();
    final ImmutableList.Builder<Result> results = ImmutableList.builder();

    do {
        final Page page = this.scanner.scan(startKey);
        startKey = page.getLastKey();
        results.addAll(page.getResults());
    } while (startKey.isPresent());

    return results.build();
}

【问题讨论】:

  • 我不认为递归 observables 创建一个巨大的调用堆栈是正确的。 scanPage 在调用下一个 scanPage 之前返回,因此调用是顺序的,但不是嵌套的。

标签: java pagination rx-java reactive-programming observable


【解决方案1】:

JohnWowUs 的回答很棒,帮助我了解了如何有效地避免递归,但有些地方我仍然感到困惑,所以我发布了我的调整版本。

总结:

  • 各个页面以Single 的形式返回。
  • 使用Flowable 流式传输页面中包含的每个项目。这意味着我们函数的调用者不需要了解各个页面,只需收集包含的项目即可。
  • 使用BehaviorProcessor 从第一页开始,并在我们检查当前页面后获取每个后续​​页面(如果下一页可用)。
  • 关键是对processor.onNext(int) 的调用会开始下一次迭代。

此代码取决于rxjavareactive-streams

import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import java.util.Optional;
import java.util.function.Function;

import io.reactivex.Flowable;
import io.reactivex.Single;
import io.reactivex.processors.BehaviorProcessor;

public class Pagination {

    // Fetch all pages and return the items contained in those pages, using the provided page fetcher function
    public static <T> Flowable<T> fetchItems(Function<Integer, Single<Page<T>>> fetchPage) {
        // Processor issues page indices
        BehaviorProcessor<Integer> processor = BehaviorProcessor.createDefault(0);
        // When an index number is issued, fetch the corresponding page
        return processor.concatMap(index -> fetchPage.apply(index).toFlowable())
                        // when returning the page, update the processor to get the next page (or stop)
                        .doOnNext(page -> {
                            if (page.hasNext()) {
                                processor.onNext(page.getNextPageIndex());
                            } else {
                                processor.onComplete();
                            }
                        })
                        .concatMapIterable(Page::getElements);
    }

    public static void main(String[] args) {
        fetchItems(Pagination::examplePageFetcher).subscribe(System.out::println);
    }

    // A function to fetch a page of our paged data
    private static Single<Page<String>> examplePageFetcher(int index) {
        return Single.just(pages.get(index));
    }

    // Create some paged data
    private static ArrayList<Page<String>> pages = new ArrayList<>(3);

    static {
        pages.add(new Page<>(Arrays.asList("one", "two"), Optional.of(1)));
        pages.add(new Page<>(Arrays.asList("three", "four"), Optional.of(2)));
        pages.add(new Page<>(Arrays.asList("five"), Optional.empty()));
    }

    static class Page<T> {
        private List<T> elements;
        private Optional<Integer> nextPageIndex;

        public Page(List<T> elements, Optional<Integer> nextPageIndex) {
            this.elements = elements;
            this.nextPageIndex = nextPageIndex;
        }

        public List<T> getElements() {
            return elements;
        }

        public int getNextPageIndex() {
            return nextPageIndex.get();
        }

        public boolean hasNext() {
            return nextPageIndex.isPresent();
        }
    }
}

输出:

one
two
three
four
five

【讨论】:

【解决方案2】:

这不是最优雅的解决方案,但您可以使用主题和副作用。请参阅下面的玩具示例

import rx.Observable;
import rx.Subscriber;
import java.util.ArrayList;
import java.util.List;
import java.util.HashMap;
import rx.subjects.*;

public class Pagination {
    static HashMap<String,ArrayList<String>> pages = new HashMap<String,ArrayList<String>>();

    public static void main(String[] args) throws InterruptedException {
        pages.put("default", new ArrayList<String>());
        pages.put("2", new ArrayList<String>());
        pages.put("3", new ArrayList<String>());
        pages.put("4", new ArrayList<String>());

        pages.get("default").add("2");
        pages.get("default").add("Maths");
        pages.get("default").add("Chemistry");  

        pages.get("2").add("3");
        pages.get("2").add("Physics");   
        pages.get("2").add("Biology"); 

        pages.get("3").add("4");
        pages.get("3").add("Art");   

        pages.get("4").add("");
        pages.get("4").add("Geography"); 



        Observable<List<String>> ret = Observable.defer(() -> 
        { 
            System.out.println("Building Observable");
            ReplaySubject<String> pagecontrol = ReplaySubject.<String>create(1);
            Observable<List<String>> ret2 = pagecontrol.asObservable().concatMap(aKey -> 
            {
                if (!aKey.equals("")) {
                    return Observable.just(pages.get(aKey)).doOnNext(page -> pagecontrol.onNext(page.get(0)));
                } else {
                    return Observable.<List<String>>empty().doOnCompleted(()->pagecontrol.onCompleted());
                }
            });
            pagecontrol.onNext("default");
            return ret2;
        });
        // Use this if you want to ensure work isn't done again
        ret = ret.cache();
        ret.subscribe(l -> System.out.println("Sub 1 : " + l));
        ret.subscribe(l -> System.out.println("Sub 2 : " + l));
        Thread.sleep(2000L);
    }
}

进行了改进。

【讨论】:

  • 抱歉,我认为这行不通。结果将只是扫描的第一个结果,因为它已被设置,但主执行线程的流程会发现 currentKey 尚未设置,然后退出。
  • 看来你是对的。问题是关键。如果没有常规的生成方式,例如“1”,“2”等,那么我认为如果您有太多页面,您会被递归解决方案所困扰。
  • 查看使用主题和副作用的 hacky 解决方案的编辑答案。
  • 这肯定更接近。我遇到的问题是我必须订阅可观察的 before 我启动它。这使得从函数返回 Observable 变得很困难。我想诀窍是传递订阅者,而不是返回流。不幸的是,这不是我的代码在 ATM 上的工作方式。
  • 经过改进进行了编辑。我认为这应该满足您的所有要求。有趣的问题。
【解决方案3】:

另一种方法是使用令牌流:获取初始令牌的数据,一旦获得新的远程数据就将下一个令牌推送到流,然后重新订阅直到令牌为空

 public Observable<Window> paging() {

        Subject<Token, Token> tokenStream = BehaviorSubject.<Token>create().toSerialized();

        tokenStream.onNext(Token.startToken());

        Observable<Window> dataStream =
                Observable.defer(() -> tokenStream.first().flatMap(this::remoteData))
                        .doOnNext(window -> tokenStream.onNext(window.getToken()))
                        .repeatWhen(completed -> completed.flatMap(__ -> tokenStream).takeWhile(Token::hasMore));

        return dataStream;
    }

结果是

Window{next token=Token{key='1'}, data='data for token: Token{key=''}'}
Window{next token=Token{key='2'}, data='data for token: Token{key='1'}'}
Window{next token=Token{key='3'}, data='data for token: Token{key='2'}'}
Window{next token=Token{key='4'}, data='data for token: Token{key='3'}'}
Window{next token=Token{key='5'}, data='data for token: Token{key='4'}'}
Window{next token=Token{key='6'}, data='data for token: Token{key='5'}'}
Window{next token=Token{key='7'}, data='data for token: Token{key='6'}'}
Window{next token=Token{key='8'}, data='data for token: Token{key='7'}'}
Window{next token=Token{key='9'}, data='data for token: Token{key='8'}'}
Window{next token=Token{key='10'}, data='data for token: Token{key='9'}'}

复制粘贴样本

public class RxPaging {

    public Observable<Window> paging() {

        Subject<Token, Token> tokenStream = BehaviorSubject.<Token>create().toSerialized();

        tokenStream.onNext(Token.startToken());

        Observable<Window> dataStream =
                Observable.defer(() -> tokenStream.first().flatMap(this::remoteData))
                        .doOnNext(window -> tokenStream.onNext(window.getToken()))
                        .repeatWhen(completed -> completed.flatMap(__ -> tokenStream).takeWhile(Token::hasMore));

        return dataStream;
    }

    private Observable<Window> remoteData(Token token) {
        /*limit number of pages*/
        int page = page(token);
        Token nextToken = page < 10
                ? nextPageToken(token)
                : Token.endToken();

        return Observable
                .just(new Window(nextToken, "data for token: " + token))
                .delay(100, TimeUnit.MILLISECONDS);
    }

    private int page(Token token) {
        String key = token.getKey();
        return key.isEmpty() ? 0 : Integer.parseInt(key);
    }

    private Token nextPageToken(Token token) {
        String tokenKey = token.getKey();
        return tokenKey.isEmpty() ? new Token("1") : nextToken(tokenKey);
    }

    private Token nextToken(String tokenKey) {
        return new Token(String.valueOf(Integer.parseInt(tokenKey) + 1));
    }

    public static class Token {
        private final String key;

        private Token(String key) {
            this.key = key;
        }

        public static Token endToken() {
            return startToken();
        }

        public static Token startToken() {
            return new Token("");
        }

        public String getKey() {
            return key;
        }

        public boolean hasMore() {
            return !key.isEmpty();
        }

        @Override
        public String toString() {
            return "Token{" +
                    "key='" + key + '\'' +
                    '}';
        }
    }


    public static class Window {
        private final Token token;
        private final String data;

        public Window(Token token, String data) {
            this.token = token;
            this.data = data;
        }

        public Token getToken() {
            return token;
        }

        public String getData() {
            return data;
        }

        @Override
        public String toString() {
            return "Window{" +
                    "next token=" + token +
                    ", data='" + data + '\'' +
                    '}';
        }
    }

    @Test
    public void testPaging() throws Exception {
        paging().toBlocking().subscribe(System.out::println);
    }
}

【讨论】:

  • 我喜欢这个。我还没有尝试实现它,但我想知道它与在 JohnWowUs 的解决方案中使用主题相比如何?
  • 两者都使用辅助令牌流 - 不同之处在于数据流“循环”如何被中止。关于主题 - 在约翰的解决方案中,使用了无限重放主题,这可能是也可能不是问题,取决于估计的页数。它可以安全地被 BehaviorSubject 替换(我认为他打错了字并使用了 ReplaySubject.create(1) 而不是 ReplaySubject.createWithSize(1) - 这与 BehaviorSubject 相同)
【解决方案4】:

只是一个想法,您为什么不实现自己的迭代器来迭代您的页面,然后从中创建观察者?

例子:

Observable.from(new Iterable<T>() {

        @Override
        public Iterator<T> iterator() {
            return new Iterator<T>() {
                @Override
                public boolean hasNext() {
                    return hasNextPage(currentPageKey);
                }

                @Override
                public T next() {

                    page = getNextPage(currentPageKey);
                    currentPageKey = page.getKey();
                    return page;
                }

                @Override
                public void remove() {
                    throw new UnsupportedOperationException();
                }
            };
        }
    });

更优雅的方法是让您的页面管理器(我相信您的代码示例中的扫描仪变量)实现可迭代并在那里编写迭代逻辑。

【讨论】:

  • 这不起作用,因为您只是实现了一个同步的 Iterable 并将其包装在一个 Observable 中。事实上,问题在于“getNextPage”返回一个 Observable,因此迭代器永远不会真正知道它是否“hasNext”。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2016-08-06
  • 1970-01-01
  • 2014-10-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多