如何收集顺序调用异步 API 的结果?

Sav*_*ior 6 java asynchronous java-8 completable-future

我有一个异步 API,它本质上是通过分页返回结果

public CompletableFuture<Response> getNext(int startFrom);
Run Code Online (Sandbox Code Playgroud)

每个Response对象都包含一个偏移量列表startFrom和一个标志,该标志指示是否还有更多剩余元素,因此是否getNext()需要发出另一个请求。

我想编写一个方法来遍历所有页面并检索所有偏移量。我可以像这样以同步方式编写它

int startFrom = 0;
List<Integer> offsets = new ArrayList<>();

for (;;) {
    CompletableFuture<Response> future = getNext(startFrom);
    Response response = future.get(); // an exception stops everything
    if (response.getOffsets().isEmpty()) {
        break; // we're done
    }
    offsets.addAll(response.getOffsets());
    if (!response.hasMore()) {
        break; // we're done
    }
    startFrom = getLast(response.getOffsets());
}
Run Code Online (Sandbox Code Playgroud)

换句话说,我们在0处调用getNext()with。startFrom如果抛出异常,我们就会短路整个过程。否则,如果没有偏移,我们就完成。如果有偏移,我们将它们添加到主列表中。如果没有更多的东西需要获取,我们就完成了。否则,我们将重置startFrom为我们获取的最后一个偏移量并重复。

理想情况下,我想在不阻塞CompletableFuture::get()并返回CompletableFuture<List<Integer>>包含所有偏移量的情况下执行此操作。

我怎样才能做到这一点?我如何编写期货来收集结果?


我正在考虑“递归”(实际上不是在执行中,而是在代码中)

private CompletableFuture<List<Integer>> recur(int startFrom, List<Integer> offsets) {
    CompletableFuture<Response> future = getNext(startFrom);
    return future.thenCompose((response) -> {
        if (response.getOffsets().isEmpty()) {
            return CompletableFuture.completedFuture(offsets);
        }
        offsets.addAll(response.getOffsets());
        if (!response.hasMore()) {
            return CompletableFuture.completedFuture(offsets);
        }
        return recur(getLast(response.getOffsets()), offsets);
    });
}

public CompletableFuture<List<Integer>> getAll() {
    List<Integer> offsets = new ArrayList<>();
    return recur(0, offsets);
}
Run Code Online (Sandbox Code Playgroud)

从复杂性的角度来看,我不喜欢这个。我们可以做得更好吗?

Did*_*r L 1

为了练习,我制作了该算法的通用版本,但它相当复杂,因为您需要:

  1. 调用服务的初始值(startFrom)
  2. 服务调用本身 ( getNext())
  3. 用于累积中间值的结果容器(offsets)
  4. 累加器 ( offsets.addAll(response.getOffsets()))
  5. 执行“递归”的条件 ( response.hasMore())
  6. 计算下一个输入的函数 ( getLast(response.getOffsets()))

所以这给出:

public <T, I, R> CompletableFuture<R> recur(T initialInput, R resultContainer,
        Function<T, CompletableFuture<I>> service,
        BiConsumer<R, I> accumulator,
        Predicate<I> continueRecursion,
        Function<I, T> nextInput) {
    return service.apply(initialInput)
            .thenCompose(response -> {
                accumulator.accept(resultContainer, response);
                if (continueRecursion.test(response)) {
                    return recur(nextInput.apply(response),
                            resultContainer, service, accumulator,
                            continueRecursion, nextInput);
                } else {
                    return CompletableFuture.completedFuture(resultContainer);
                }
            });
}

public CompletableFuture<List<Integer>> getAll() {
    return recur(0, new ArrayList<>(), this::getNext,
            (list, response) -> list.addAll(response.getOffsets()),
            Response::hasMore,
            r -> getLast(r.getOffsets()));
}
Run Code Online (Sandbox Code Playgroud)

recur()可以通过用第一次调用的结果返回的值替换initialInput来进行小的简化,并且可以将 和 合并为一个,然后可以将 与该函数合并。CompletableFutureresultContaineraccumulatorConsumerservicenextInput

但这会变得更复杂一些getAll():

private <I> CompletableFuture<Void> recur(CompletableFuture<I> future,
        Consumer<I> accumulator,
        Predicate<I> continueRecursion,
        Function<I, CompletableFuture<I>> service) {
    return future.thenCompose(result -> {
        accumulator.accept(result);
        if (continueRecursion.test(result)) {
            return recur(service.apply(result), accumulator, continueRecursion, service);
        } else {
            return CompletableFuture.completedFuture(null);
        }
    });
}

public CompletableFuture<List<Integer>> getAll() {
    ArrayList<Integer> resultContainer = new ArrayList<>();
    return recur(getNext(0),
            result -> resultContainer.addAll(result.getOffsets()),
            Response::hasMore,
            r -> getNext(getLast(r.getOffsets())))
            .thenApply(unused -> resultContainer);
}
Run Code Online (Sandbox Code Playgroud)