【发布时间】: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());
wrapPredicate 和 wrapFunction 仅用于检查异常重新抛出。
但是,显然,对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.allOf将List<CompletableFuture<Optional<T>>>合并为CompletableFuture<List<Optional<T>>>,我最终只等待一个未来。
标签: java java-stream