通过拆分和运行将ListenableFuture <Iterable <A >>转换为Iterable <ListenableFuture <B >>

Ben*_*ith 7 java concurrency asynchronous future guava

我正在寻找将一个人转换ListenableFuture<Iterable<A>>成一系列个人的最好方法ListenableFutures.这是我正在寻找的那种方法签名:

public <A, B> Iterable<ListenableFuture<B>> splitAndRun(
    final ListenableFuture<Iterable<A>> elements, 
    final Function<A, B> func, 
    final ListeningExecutorService executor
);
Run Code Online (Sandbox Code Playgroud)

显然,如果我回来ListenableFuture<Iterable<ListenableFuture<B>>>,我可以做到,但我觉得我应该能够分裂并运行它并保持其异步性.

这是我到目前为止的代码,但你会注意到最后的讨厌.get(),这会破坏事物.如果我的事情过于复杂,请原谅.

public class CallableFunction<I, O> implements Callable<O>{
  private final I input;
  private final Function<I, O> func;

  public CallableFunction(I input, Function<I, O> func) {
    this.input = input;
    this.func = func;
  }

  @Override public O call() throws Exception {
    return func.apply(input);
  }
}

public <A, B> Iterable<ListenableFuture<B>> splitAndRun(
    final ListenableFuture<Iterable<A>> elements, 
    final Function<A, B> func, 
    final ListeningExecutorService executor
) throws InterruptedException, ExecutionException {
  return Futures.transform(elements, 
      new Function<Iterable<A>, Iterable<ListenableFuture<B>>>() {
    @Override
    public Iterable<ListenableFuture<B>> apply(Iterable<A> input) {
      return Iterables.transform(input, new Function<A, ListenableFuture<B>>() {
        @Override
        public ListenableFuture<B> apply(A a) {
          return executor.submit(new CallableFunction<A, B>(a, func));
        }
      });
    }
  }, executor).get();
}
Run Code Online (Sandbox Code Playgroud)

Chr*_*irk 2

(替代我原来的答案

但是,如果转型速度缓慢,或者某些投入可能会失败,但其他投入可能会成功,该怎么办?在这种情况下,我们希望单独转换每个输出。另外,我们希望确保转换只发生一次。我们的集合转换方法此保证。因此,在您的示例代码中,输出的每次迭代都会向执行器提交新任务,即使之前提交的任务可能已经完成。出于这个原因,我们Iterables.transform只在改造轻量级时才推荐和朋友。(通常,如果您正在做一些重量级的事情,您的转换函数将抛出一个已检查的异常,这是Function不允许的。将此视为提示:)当然,您的示例恰好不会触发提示。)

这一切在代码中意味着什么?基本上,我们将颠倒我在其他答案中给出的操作顺序。我们将首先从 转换Future<Iterable<A>>Iterable<Future<A>>然后我们将为每个任务提交一个单独的任务A,将其转换为B. 对于后一步,我们将提供一个Executor ,以便转换不会阻塞一些无辜的线程。(我们现在只需要Futures.transform,所以我已经静态导入了它。)

List<ListenableFuture<A>> individuals = newArrayList();
for (int i = 0; i < knownSize; i++) {
  final int index = i;
  individuals.add(transform(input, new Function<List<A>, A>() {
    @Override
    public A apply(List<A> values) {
      return values.get(index);
    }
  }));
}

List<ListenableFuture<B>> result = newArrayList();
for (ListenableFuture<A> original : individuals) {
  result.add(transform(original, function, executor));
}
return result;
Run Code Online (Sandbox Code Playgroud)

无论如何,这就是这个想法。但我的实现很愚蠢。我们可以轻松地同时执行这两个步骤:

List<ListenableFuture<B>> result = newArrayList();
for (int i = 0; i < knownSize; i++) {
  final int index = i;
  result.add(transform(input, new Function<List<A>, B>() {
    @Override
    public B apply(List<A> values) {
      return function.apply(values.get(index));
    }
  }, executor));
}
return result;
Run Code Online (Sandbox Code Playgroud)

因为这会进行 nFutures.transform次调用而不是 1 次,并且因为它使用单独的Executor,所以如果转换是重量级的,那么它比我的其他解决方案更好,如果转换是轻量级的,那么它会比我的其他解决方案更好。另一个警告仍然存在:只有当您知道将有多少输出时,这才有效。