【问题标题】:New to Spark, mapping with graphx graphs - NullPointerExceptionSpark 的新手,使用 graphx 图形进行映射 - NullPointerException
【发布时间】:2016-09-28 17:23:37
【问题描述】:

我的目标是从一个共同的完整图中计算多个子图中的三角形。子图由一组常量节点 + 来自 RDD[Long] 的节点定义。我是 spark/graphx 的新手,所以这可能是对 map 的不当使用。 我分享的代码会重现我的错误。

首先,我有一个完整图的子图,如下所示

import org.apache.spark.rdd._
import org.apache.spark.graphx._
val nodes: RDD[(VertexId, String)] = sc.parallelize(Array((3L, "3"), (7L, "7"), (5L, "5"), (2L, "2"),(4L,"4")))
val vertices: RDD[Edge[String]] = sc.parallelize(Array(Edge(3L, 7L, "a"), Edge(3L, 5L, "b"), Edge(2L, 5L, "c"), Edge(5L, 7L, "d"), Edge(2L, 7L, "e"),Edge(4L,5L,"f")))
val graph: Graph[String,String] = Graph(nodes, vertices, "z")

val baseNodes: Array[Long] = Array(2L,5L,7L)    
val subgraph = graph.subgraph(vpred = (vid,attr)=> baseNodes contains vid)

然后我从图中声明其他节点的 RDD[Long]。

val testNodes: RDD[Long] = sc.parallelize(Array(3L,4L))

我想将每个 testNode 添加到子图中并计算 testNode 上存在的三角形。

val triangles: RDD[(Long,Int)] = testNodes.map{ newNode =>
  val newNodes: Array[Long] = baseNodes :+ newNode
  val newSubgraph = graph.subgraph(vpred = (vid,attr)=> newNodes contains vid)
  (newNode,findTriangles(7L,newSubgraph))
}
triangles.foreach(x=>x.toString)

如果我在 map 函数之外调用它,我的 findTriangles 可以正常工作。

def findTriangles(id:Long,subgraph:Graph[String,String]): Int = {
  val triCounts = subgraph.triangleCount().vertices
  val count:Int = triCounts.filter{case(item,count)=> {item.toInt == id}}.map{case(item,count)=>count}.first
  count
}
val triangles = findTriangles(7L,subgraph) //1

但是当我运行我的地图函数来计算三角形时,我得到了 NullPointerException。我认为问题在于在映射函数中使用我的图形验证。是这个问题吗?有没有办法解决这个问题?

【问题讨论】:

  • 您不能在另一个RDD 中处理'RDD。这不仅限于GraphX——它是对RDDs 的限制,因为GraphX 是建立在RDDs 之上的,嗯——那是你的问题。长话短说——执行者对你的其他RDDs 一无所知,因此NPE。一般来说,解决这个问题的方法是加入——我没有深入研究你想要做什么,但这是一般的想法。

标签: scala apache-spark rdd spark-graphx map-function


【解决方案1】:

我认为问题应该是 baseNodes 变量。本地声明的变量,例如示例中的 baseNode,仅在 Spark 驱动程序中可见,而不在实际执行转换和操作的执行程序中可见。为了避免 NullPointerException,您需要并行化在执行程序上执行的转换(如映射)中需要的任何变量。作为替代方案,如果您拥有的变量是只读的,您可以使用 Spark 中的广播构造将该变量广播给执行器。在您的情况下,似乎 baseNodes 没有在 map 操作中被修改,因此它是广播而不是并行化的一个很好的候选者。

【讨论】:

  • 感谢您的回复!在您的回答和大卫格里芬的评论之间,我能够找出我的问题。广播解决了非 RDD 变量的问题,但由于 GraphX 中的图是建立在 RDD 之上的,因此无法在另一个 map 函数中访问它们。似乎 graphx 是为在内存中的单个大图上扩展而构建的,而不是将同一图的副本传播到多个执行器。由于我的图表不是非常大,我放弃了 graphx 并将图表存储为数组和地图并实现了我自己的 triangleCount,让我可以广播图表
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2015-10-26
  • 1970-01-01
  • 1970-01-01
  • 2022-06-10
  • 2020-06-13
  • 1970-01-01
  • 2016-01-06
相关资源
最近更新 更多