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)
从复杂性的角度来看,我不喜欢这个。我们可以做得更好吗?
为了练习,我制作了该算法的通用版本,但它相当复杂,因为您需要:
startFrom)getNext())offsets)offsets.addAll(response.getOffsets()))response.hasMore())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)
| 归档时间: |
|
| 查看次数: |
1772 次 |
| 最近记录: |