【问题标题】:How to run a graph algorithm concurrently in Java using multi-core parallelism如何使用多核并行在 Java 中同时运行图形算法
【发布时间】:2017-11-09 04:02:45
【问题描述】:

我想在大型图上同时运行一个算法,使用多核并行。我一直在研究它一段时间,但一直未能提出一个好的解决方案。

这是简单的算法:

W - a very large number
double weight = 0

while(weight < W)

    - v : get_random_node_from(Graph)

    - weight += calculate(v)
  • 我研究了 fork-and-join,但找不到将这个问题分成更小的子问题的方法。
  • 然后我尝试使用 Java 8 流,为此我需要创建一个 lambda 表达式。当我尝试做这样的事情时:

double weight = 0 Callable<Object> task = () -> { can not update weight here, as it needs to be final }

我的问题是,是否可以在 lambda 方法中更新像 weight 这样的变量?或者有没有更好的方法可以解决这个问题?

我得到的最接近的是使用ExecutorService,但遇到了同步问题。

------------编辑-------------

这里是详细的算法:

简而言之,我要做的就是遍历一个海量图,对随机选择的节点执行操作(只要权重

这需要的时间太长,因为它没有利用 CPU 的全部功能。

理想情况下,多核上的所有线程/进程将在随机选择的节点上执行操作,并更新共享的 weightIndex

注意:不同线程是否选择同一个节点无关紧要,因为它是随机的,没有替换。

算法:

函数串行(){

List<List<Integer>> I (shared data structure which I want to update)
double weight

//// Task which I want to parallelize

while(weight < W) {

    v : get_random_node_from(Graph)

    bfs(v, affected_nodes) ...// this will fill up affected_nodes by v

    foreach(affected_node in affected_nodes) {

         // update I related to affected_node
         // and do other computation
    }

    weight += affected_nodes.size()

}

///////// Parallelization ends here

use_index(I) // I is passed now to some other method(not important) to get further results

}

重要的是,所有线程都更新相同的Iweight

谢谢。

【问题讨论】:

  • 我不确定您要完成什么 - 如果要计算整个图的权重,使用 Java 流,您只需让流的每个元素成为图表,使用 map 将其映射到权重,然后对流求和。这也应该是可并行化的。
  • weight &lt; W。你开始了几个calculate(v),但是在第一个weight += calculate(v)之后weight变成&gt;=W。是否应该取消其他计算,或者他们可以完成他们的工作?如果他们可以完成他们的工作,他们可以并且应该将他们的结果添加到weight
  • @AlexeiKaigorodov 每当weight &gt; W 所有进程都应该停止。我将编辑问题以添加更多详细信息。
  • Callable 传递给 executor 不必是 la​​mbda。它可以是任何实现Callable 的类型,因此可以访问任何信息。
  • 感谢您的建议。但是,我已经尝试过了,但遇到了同步问题(可能)并得到了 invalidOperationExceptions。

标签: concurrency java-8 executorservice multicore java.util.concurrent


【解决方案1】:

好吧,你可以将 weight 包装到一个包含单个元素的数组中,这对于这类东西来说是一种诀窍;甚至由 java 在内部完成,如下所示:

weight[0] = weight[0] + calculate(v);

但是这样做存在问题,因为您要并行运行它。你不会得到你想要的结果,因为weight[0] 不是线程安全的。您可以使用某种同步,但 java 已经有一个很好的解决方案:DoubleAdder 在竞争环境(和多个 cpu)中可以更好地扩展。

一个微不足道的小例子:

DoubleAdder weight = new DoubleAdder();

private static int calculate(int v) {
    return v + 1;
}


Stream.of(1, 2, 3, 4, 5, 6, 7, 8, 9)
            .parallel()
            .forEach(x -> {
                int y = calculate(x);
                weight.add(y);
            });

System.out.println(weight); // 54

然后是您要为此选择的随机化器的问题:get_random_node_from(Graph)。您确实需要获得一个随机的Node,但同时您需要只获得一次。 但是,如果您可以将所有节点 flatten 合并到一个 List 中,那么您可能不需要它。

这里的问题是Graphs通常以递归方式遍历,你不知道它的确切大小:

while(parent.hasChildren) {
     traverse children and so on...
}

这会在Streams下并行化不好,你可以看看Spliterators#spliteratorUnknownSize。它将从1024 算术增长;这就是为什么我建议将节点扁平化为一个已知大小的列表;这将更好地并行化。

【讨论】:

    猜你喜欢
    • 2021-10-03
    • 1970-01-01
    • 1970-01-01
    • 2015-12-12
    • 1970-01-01
    • 2017-05-13
    • 1970-01-01
    • 1970-01-01
    • 2016-03-03
    相关资源
    最近更新 更多