【问题标题】:Spark - GraphX - scaling connected componentsSpark - GraphX - 缩放连接的组件
【发布时间】:2017-03-09 01:57:19
【问题描述】:

我正在尝试使用连接的组件,但在缩放时遇到了问题。我的这就是我所拥有的-

// get vertices
val vertices = stage_2.flatMap(x => GraphUtil.getVertices(x)).cache

// get edges
val edges = stage_2.map(x => GraphUtil.getEdges(x)).filter(_ != null).flatMap(x => x).cache

// create graph  
val identityGraph = Graph(vertices, edges)

// get connected components
val cc = identityGraph.connectedComponents.vertices

其中,GraphUtil 具有返回顶点和边的辅助函数。此时,我的图有约 100 万个节点和约 200 万条边(顺便说一句,预计将增长到约 1 亿个节点)。我的图连接非常稀疏 - 所以我希望有很多小图。

当我运行上述代码时,我不断收到java.lang.OutOfMemoryError: Java heap space。我尝试使用executor-memory 32g 并运行一个由 15 个节点组成的集群,其中 45g 作为纱线容器大小。

这里是异常详情:

16/10/26 10:32:26 ERROR util.Utils: uncaught error in thread SparkListenerBus, stopping SparkContext
java.lang.OutOfMemoryError: Java heap space
    at java.util.Arrays.copyOfRange(Arrays.java:2694)
    at java.lang.String.<init>(String.java:203)
    at java.lang.StringBuilder.toString(StringBuilder.java:405)
    at com.fasterxml.jackson.core.util.TextBuffer.contentsAsString(TextBuffer.java:360)
    at com.fasterxml.jackson.core.io.SegmentedStringWriter.getAndClear(SegmentedStringWriter.java:98)
    at com.fasterxml.jackson.databind.ObjectMapper.writeValueAsString(ObjectMapper.java:2216)
    at org.json4s.jackson.JsonMethods$class.compact(JsonMethods.scala:32)
    at org.json4s.jackson.JsonMethods$.compact(JsonMethods.scala:44)
    at org.apache.spark.scheduler.EventLoggingListener$$anonfun$logEvent$1.apply(EventLoggingListener.scala:146)
    at org.apache.spark.scheduler.EventLoggingListener$$anonfun$logEvent$1.apply(EventLoggingListener.scala:146)
    at scala.Option.foreach(Option.scala:236)
    at org.apache.spark.scheduler.EventLoggingListener.logEvent(EventLoggingListener.scala:146)
    at org.apache.spark.scheduler.EventLoggingListener.onJobStart(EventLoggingListener.scala:173)
    at org.apache.spark.scheduler.SparkListenerBus$class.onPostEvent(SparkListenerBus.scala:34)
    at org.apache.spark.scheduler.LiveListenerBus.onPostEvent(LiveListenerBus.scala:31)
    at org.apache.spark.scheduler.LiveListenerBus.onPostEvent(LiveListenerBus.scala:31)
    at org.apache.spark.util.ListenerBus$class.postToAll(ListenerBus.scala:55)
    at org.apache.spark.util.AsynchronousListenerBus.postToAll(AsynchronousListenerBus.scala:37)
    at org.apache.spark.util.AsynchronousListenerBus$$anon$1$$anonfun$run$1$$anonfun$apply$mcV$sp$1.apply$mcV$sp(AsynchronousListenerBus.scala:80)
    at org.apache.spark.util.AsynchronousListenerBus$$anon$1$$anonfun$run$1$$anonfun$apply$mcV$sp$1.apply(AsynchronousListenerBus.scala:65)
    at org.apache.spark.util.AsynchronousListenerBus$$anon$1$$anonfun$run$1$$anonfun$apply$mcV$sp$1.apply(AsynchronousListenerBus.scala:65)
    at scala.util.DynamicVariable.withValue(DynamicVariable.scala:57)
    at org.apache.spark.util.AsynchronousListenerBus$$anon$1$$anonfun$run$1.apply$mcV$sp(AsynchronousListenerBus.scala:64)
    at org.apache.spark.util.Utils$.tryOrStopSparkContext(Utils.scala:1181)
    at org.apache.spark.util.AsynchronousListenerBus$$anon$1.run(AsynchronousListenerBus.scala:63)

此外,我收到了大量以下日志:

16/10/26 10:30:32 INFO spark.MapOutputTrackerMaster: Size of output statuses for shuffle 320 is 263 bytes
16/10/26 10:30:32 INFO spark.MapOutputTrackerMaster: Size of output statuses for shuffle 321 is 268 bytes
16/10/26 10:30:32 INFO spark.MapOutputTrackerMaster: Size of output statuses for shuffle 322 is 264 bytes

我的问题是有人尝试过这种规模的 ConnectedComponents 吗?如果是,我做错了什么?

【问题讨论】:

    标签: apache-spark spark-graphx connected-components


    【解决方案1】:

    正如我在上面的 cmets 中发布的,我在 Spark 上使用 map/reduce 实现了连接组件。您可以在这里找到更多详细信息 - https://www.linkedin.com/pulse/connected-component-using-map-reduce-apache-spark-shirish-kumar 和 MIT 许可下的源代码 - https://github.com/kwartile/connected-component

    【讨论】:

    • 这与原生 spark 连接组件的性能/可扩展性相比如何?
    • 我假设您所说的原生是指 GraphX 实现。上次使用 GraphX(大约是一年前),它对我们来说并没有扩展。到目前为止,在我们的测试中,我们的实现运行良好。我已将其记录为自述文件。
    【解决方案2】:

    连通分量算法的扩展性不是很好,它的性能很大程度上取决于你的图的拓扑结构。你的边缘稀疏并不意味着你有小组件。一长串边非常稀疏(边数 = 顶点数 - 1),但在 GraphX 中实现的蛮力算法效率不会很高(参见 ccpregel 的来源)。

    这是您可以尝试的(已排序,仅代码):

    1. 检查镶木地板中的顶点和边(在磁盘上),然后再次加载它们以构建图形。当您的执行计划变得太大时,缓存有时并不能解决问题。
    2. 以保持算法结果不变的方式转换图形。例如,您可以在code 中看到算法正在双向传播信息(默认情况下应该如此)。因此,如果您有多个连接相同两个顶点的边,请将它们从应用算法的图形中过滤掉。
    3. 自己优化 GraphX 代码(这真的很简单),使用通用优化节省内存(即在每次迭代时在磁盘上设置检查点以避免 OOM),或特定领域的优化(类似于第 2 点)

    如果您可以将 GraphX(它变得有些遗留)抛在脑后,您可以考虑使用 GraphFrames(packageblog )。没试过,不知道有没有CC。

    我确信您可以在 spark 包中找到其他可能性,但也许您甚至想使用 Spark 之外的东西。但这超出了问题的范围。

    祝你好运!

    【讨论】:

    • GraphFrames 在后台使用 DataFrames 和 GraphX,所以我不明白这将如何帮助 OP。
    • @eliasah 我希望 GraphFrames 会比 GraphX 更优化。它们在 DataFrames 上运行的事实是一个好兆头,因为这样它们就可以利用催化剂优化器和钨。正如我所说,我没有尝试,我只是抱着一个明智的希望。
    • GraphX 是一个没人愿意从事的项目,因为底层的图论很难扩展。不幸的是,它“几乎”是一个死项目。尽管如此,我认为 GraphFrames 不会比 GraphX 走得更远。
    • 那么,你必须阅读这篇文章apache-spark-developers-list.1001551.n3.nabble.com/…spark 周围发生了很多我不太高兴的事情,我认为自己是 spark 狂热者。
    • 感谢您的回复。我决定在使用 map/reduce 构建的 ConnectedComponent 上从 GraphX ConnectedComponent 切换到我自己的版本 - 到目前为止表现还不错。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-04-15
    • 2016-01-25
    相关资源
    最近更新 更多