【发布时间】: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