【发布时间】:2020-08-11 06:21:44
【问题描述】:
我在 Spark GraphX 中使用 Pregel 编写了我的算法。但不幸的是,我得到了 TypeMismatch 错误。
我加载图表:val my_graph= GraphLoader.edgeListFile(sc, path)。所以开始的节点有这样的结构:
(1,1)
(2,1)
(3,1)
以nodeID为key,1为默认属性。
在run2 函数中,首先我更改了结构,以使每个节点都可以存储多个属性。因为我正在研究重叠社区检测算法,所以属性是标签和它们的分数。
在run2第一次运行时,每个节点的结构如下:
(34,Map(34 -> (1.0,34)))
(13,Map(13 -> (1.0,13)))
(4,Map(4 -> (1.0,4)))
(16,Map(16 -> (1.0,16)))
(22,Map(22 -> (1.0,22)))
这意味着节点 34,标签为 34,其分数等于 1。然后每个节点可以存储从其邻居接收的多个属性,并在接下来的步骤中将它们发送给其邻居。
在算法结束时,每个节点可以包含多个属性或仅包含一个属性,例如以下结构:
(1,Map((2->(0.49,1),(8->(0.9,1)),(13->(0.79,1))))
(2,Map((11->(0.89,2)),(6->(0.68,2)),(13->(0.79,2)),(10->(0.57,2))))
(3,Map((20->(0.0.8,3)),(1->(0.66,3))))
如上图,节点1属于社区2,得分0.49,属于社区8,得分0.9,属于社区13,得分0.79。
以下代码显示了 Pregel 中定义的不同函数。
def run2[VD, ED: ClassTag](graph: Graph[VD, ED], maxSteps: Int) = {
val temp_graph = graph.mapVertices { case (vid, _) => mutable.HashMap[VertexId, (Double,VertexId)](vid -> (1,vid)) }
def sendMessage(e: EdgeTriplet[mutable.HashMap[VertexId, (Double,VertexId)], ED]): Iterator[(VertexId, mutable.HashMap[VertexId, (Double, VertexId)])] = {
Iterator((e.srcId,e.dstAttr), (e.dstId,e.srcAttr))
}
def mergeMessage(count1: (mutable.HashMap[VertexId, (Double,VertexId)]), count2: (mutable.HashMap[VertexId, (Double,VertexId)]))= {
val communityMap = new mutable.HashMap[VertexId, List[(Double, VertexId)]]
(count1.keySet ++ count2.keySet).map(key => {
val count1Val = count1.getOrElse(key, (0D,0:VertexId))
val count2Val = count2.getOrElse(key, (0D,0:VertexId))
communityMap += key->(count1Val::communityMap(key))
communityMap += key->(count2Val::communityMap(key))
})
communityMap
}
def vertexProgram(vid: VertexId, attr: mutable.HashMap[VertexId,(Double, VertexId)], message: mutable.HashMap[VertexId, List[(Double, VertexId)]]) = {
if (message.isEmpty)
attr
else {
val labels_score: mutable.HashMap[VertexId, Double] = message.map {
key =>
var value_sum = 0D
var isMemberFlag = 0
var maxSimilar_result = 0D
val max_similar = most_similar.filter(x=>x._1==vid)(1)
if (key._2.exists(x=>x._2==max_similar)) isMemberFlag = 1 else isMemberFlag = 0
key._2.map {
values =>
if (values._2==max_similar) maxSimilar_result = values._1 else maxSimilar_result = 0D
val temp = broadcastVariable.value(vid)(values._2)._2
value_sum += values._1 * temp
}
value_sum += (beta*value_sum)+((1-beta)*maxSimilar_result)
(key._1,value_sum) //label list
}
val max_value = labels_score.maxBy(x=>x._2)._2.toDouble
val dividedByMax = labels_score.map(x=>(x._1,x._2/max_value)) // divide by maximum value
val resultMap: mutable.HashMap[VertexId,Double] = new mutable.HashMap[VertexId, Double]
dividedByMax.foreach{ row => // select labels more than threshold P = 0.5
if (row._2 >= p) resultMap += row
}
val max_for_normalize= resultMap.values.sum
val res = resultMap.map(x=>(x._1->(x._2/max_for_normalize,x._1))) // Normalize labels
res
}
}
val initialMessage = mutable.HashMap[VertexId, (Double,VertexId)]()
val overlapCommunitiesGraph = Pregel(temp_graph, initialMessage, maxIterations = maxSteps)(
vprog = vertexProgram,
sendMsg = sendMessage,
mergeMsg = mergeMessage)
overlapCommunitiesGraph
}
val my_graph= GraphLoader.edgeListFile(sc, path)
val new_updated_graph2 = run2(my_graph, 1)
在上面的代码中,p=0.5 和 beta=0.5。 most_similar 是一个 RDD,包含每个节点及其最重要的节点。例如(1,3)表示节点3是节点1最相似的邻居。broadcatVariable结构如下:
(19,Map(33 -> (1.399158675718661,0.6335049099178383), 34 -> (1.4267350687130098,0.6427405501408145)))
(15,Map(33 -> (1.399158675718661,0.6335049099178383), 34 -> (1.4267350687130098,0.6427405501408145)))
...
该结构将节点之间的关系显示为键,将其邻居作为值。例如,节点 19 与节点 33 和 34 是邻居,关系通过它们之间的分数来表示。
在算法中,每个节点发送每个属性Map,其中包含多个标签及其分数。然后在mergeMessage函数中,将相同编号的标签的值放入List,并在vertexProgram中为每个标签或键处理其列表。
更新
根据下图中的等式,我使用List 来收集标签的不同分数并在vertexProgram 函数中处理它们。因为我需要P_ji来处理每个节点的标签分数,所以我不知道是否可以在mergeMessage函数中执行或者是否需要在vertexProgram中执行。 P_ji 是源节点与其邻居之间的分数,应该乘以标签分数。
我得到的错误显示在vprog = vertexProgram, 行的前面,如图所示。谁能帮我解决这个错误?
【问题讨论】:
-
您能否添加用于创建图形的代码以及所有必需变量的值(否则无法运行)。目前缺少
most_similar、p、beta、broadcastVariable的值。 -
感谢您的回复。我已经编辑了问题的文本并解释了所有必要的事情。如果你能帮助我,我将非常感激。因为这个问题已经解决了好几天,我对此感到困惑。
-
我认为问题是由于
HashMap[VertexId, List[(Double, VertexId)]]和HashMap[VertexId, (Double, VertexId)]的混合造成的。尤其是mergeMessage,它将不带列表的HashMap作为输入,并在输出中返回带有列表的HashMap。这里输入输出类型需要相同,否则合并后的消息不能再次合并。 -
感谢您的帮助。对不起,我问,但是否可以指导我可以更改代码的哪一部分来解决这个问题?我的意思是应该更改节点属性或我应该更改哪个部分。我真的很困惑
-
最简单的方法应该是不使用列表,但是否可能取决于实际逻辑。有没有什么方法可以写出
mergeMessage,而输出不是HashMap[VertexId, (Double, VertexId)]?
标签: scala apache-spark spark-graphx