【问题标题】:Force Java stream to execute part of a pipeline early to submit blocking tasks to a thread pool强制 Java 流提前执行部分管道以将阻塞任务提交到线程池
【发布时间】:2017-04-24 15:19:36
【问题描述】:

我有一个我想要处理的对象列表,Java8 流 API 看起来是最简洁易读的方式。

但是我需要对这些对象执行的一些操作包括阻塞 IO(例如读取数据库)——所以我想将这些操作提交到具有几十个线程的线程池。

一开始我想做一些类似的事情:

myObjectList
    .stream()
    .filter(wrapPredicate(obj -> threadPoolExecutor.submit(
            () -> longQuery(obj)          // returns boolean
    ).get())                              // wait for future & unwrap boolean
    .map(filtered -> threadPoolExecutor.submit(
            () -> anotherQuery(filtered)  // returns Optional
    ))
    .map(wrapFunction(Future::get))
    .filter(Optional::isPresent)
    .map(Optional::get)
    .collect(toList());

wrapPredicatewrapFunction 仅用于检查异常重新抛出。

但是,显然,对Future.get() 的调用将阻塞流的线程,直到对给定对象的查询完成,并且在此之前流不会继续。所以一次只处理一个对象,线程池没有意义。

我可以使用并行流,但我需要希望默认的 ForkJoinPool 足以满足此要求。或者只是增加"java.util.concurrent.ForkJoinPool.common.parallelism",但我不想为了那个流而改变整个应用程序的设置。我可以在自定义ForkJoinPool 中创建流,但我看到it does not guarantee that level of parallelism

所以我最终得到了类似的东西,只是为了保证在等待期货完成之前将所有需要的任务提交到线程池:

myObjectList
    .stream()
    .map(obj -> Pair.of(obj, threadPoolExecutor.submit(
                    () -> longQuery(obj)             // returns boolean
        ))
    )
    .collect(toList()).stream()                      // terminate stream to actually submit tasks to the pool
    .filter(wrapPredicate(p -> p.getRight().get()))  // wait & unwrap future after all tasks are submitted
    .map(Pair::getLeft)
    .map(filtered -> threadPoolExecutor.submit(
            () -> anotherQuery(filtered)             // returns Optional
    ))
    .collect(toList()).stream()                      // terminate stream to actually submit tasks to the pool
    .map(wrapFunction(Future::get))                  // wait & unwrap futures after all submitted
    .filter(Optional::isPresent)
    .map(Optional::get)
    .collect(toList());

有没有明显更好的方法来实现这一点?

一种更优雅的方式来告诉流“现在对流中的每个对象执行流水线步骤”,然后继续处理.collect(toList()).stream() 以外的其他处理,以及过滤Future 效果的更好方法比将其打包到 Apache Commons Pair 以便稍后在 Pair::getRight 上过滤?或者也许是一个完全不同的方法来解决这个问题?

【问题讨论】:

  • 也许只是我比流人更老派;但我不觉得这个练习的结果如此令人愉快。有点难读;并获得所有微妙的细节将使“希望我永远不必触摸和修改此代码”。因为我太害怕打破它。
  • IMO 当方法链中间没有.collect(toList()).stream() 时它是非常可读的——但是它并没有完全按照我的意愿去做。我认为如果在这里使用经典的for 循环,它不会更易读和更易于维护(但也许?我会尝试询问同事他们的想法)。
  • 好吧,可能是for循环和其他方法。
  • 几乎,但不是真的。我提交longQuery 以获得一个布尔值,我在该布尔值上filter 原始对象。然后我提交anotherQuery 和原始对象,但只提交longQuery 返回true 的对象。其余的都是正确的——然后我阻塞等待解包的Optionals 我再次收集到一个列表。
  • 好的,知道了。我不会使用布尔结果从流中过滤元素。相反,如果结果为真,我会调用anotherQuery,如果结果为假,我什么也不做。我也永远不会在流中调用Future.get。我只是将期货收集到一个列表中,然后使用 Guava 或CompletableFuture.allOfList<CompletableFuture<Optional<T>>> 合并为CompletableFuture<List<Optional<T>>>,我最终只等待一个未来。

标签: java java-stream


【解决方案1】:

您可以通过使用大大简化您的代码

myObjectList.stream()
    .map(obj -> threadPoolExecutor.submit(
                    () -> longQuery(obj)? anotherQuery(obj).orElse(null): null))
    .collect(toList()).stream()
    .map(wrapFunction(Future::get))
    .filter(Objects::nonNull)
    .collect(toList());

有一点是,如果您稍后将anotherQuery 提交给同一个执行者,并发性不会有任何改善。因此,您可以在longQuery 返回true 之后直接执行它。此时,obj 仍在范围内,因此您可以将其用于anotherQuery

通过提取Optional的结果,使用null作为失败的表示,我们可以获得相同的缺失结果表示,因为longQuery返回falseanotherQuery返回一个空Optional。所以在提取Future的结果后,我们要做的就是到.filter(Objects::nonNull)

你必须先提交作业,收集Futures,然后才能得到实际结果的逻辑没有改变。无论如何也没有办法。其他便利方法或框架所能提供的只是隐藏这些对象的实际临时存储。

【讨论】:

    【解决方案2】:

    我认为主要问题的答案是否定的。为了“执行”一个流,你需要一个终端操作。但可能还有改进的余地。

    您至少可以通过收集到地图而不是列表来摆脱这对:

    stream.collect(toMap(Function.identity(),
                         obj -> threadPoolExecutor.submit(() -> longQuery(obj))))
          .entrySet()
          .stream()
          .filter(wrapPredicate(entry -> entry.getValue().get()))
          .map(Entry::getKey)
          ...
    

    请注意,这仅在处理的对象都不等于另一个时才有效。它使代码略短且更易于阅读,因为您不必自己创建对/条目。

    【讨论】:

    • 它们是(或至少应该是)独一无二的,所以收集到地图是一个非常好的建议,我没有考虑过。谢谢。 :)
    【解决方案3】:

    您可以为 Java 8 并行流指定线程池。您不必更改应用程序设置。更多信息:https://stackoverflow.com/a/22269778/7123191

    【讨论】:

    • 正如我在描述中解释的那样——在自定义大 ForkJoinPool 中运行并行流并不能保证所需的并行化水平,因此它不是解决此问题的方法。如果我能做到这一点就完美了,但是并行流的运行批次仍然可能比我指定的线程数少。比较stackoverflow.com/a/29272776/359949
    【解决方案4】:

    您至少可以通过收集到地图而不是列表来摆脱这对:

    stream.collect(toMap(Function.identity(),
                         obj -> threadPoolExecutor.submit(() -> longQuery(obj))))
          .entrySet()
          .stream()
          .filter(wrapPredicate(entry -> entry.getValue().get()))
          .map(Entry::getKey)
          ...
    

    【讨论】:

    猜你喜欢
    • 2016-03-15
    • 1970-01-01
    • 1970-01-01
    • 2016-06-24
    • 1970-01-01
    • 2017-11-22
    • 2021-12-28
    • 2011-12-05
    相关资源
    最近更新 更多