现实生活中存在各种各样的网络,诸如人际关系网、交易网、运输网等等。对这些网络进行社区发现具有极大的意义,如在人际关系网中,可以发现出具有不同兴趣、背景的社会团体,方便进行不同的宣传策略;在交易网中,不同的社区代表不同购买力的客户群体,方便运营为他们推荐合适的商品;在资金网络中,社区有可能是潜在的洗钱团伙、刷钻联盟,方便安全部门进行相应处理;在相似店铺网络中,社区发现可以检测出商帮、价格联盟等,对商家进行指导等等。总的来看,社区发现在各种具体的网络中都能有重点的应用场景,图1展示了基于图的拓扑结构进行社区发现的例子。
如上所述,在一个无向图中,点和点之间通过一条无向边链接。但是有的时候我们需要根据图中已经给的信息来挖掘可能存在一定共性的社区。Graphx中的LPA算法就是用来实现这一功能的。
LPA算法是一个非常简单的迭代算法,通过向邻近的点发送本点的标签,在接收到标签的时候合并收到的标签,按照同一个key值合并,计算同一个key值的数量,然后用出现最多次的标签来更新本node的标签。
最后,根据每个Node上的标签来确定是否是属于同一个社区。这个算法在计算关系图中可以挖掘出比较有价值的信息。
以下是Graphx中的源码分析,在分析之前需要你了解一些关于Pregel的基础知识。如果有需要请翻阅之前的blog。
def run[VD, ED: ClassTag](graph: Graph[VD, ED], maxSteps: Int): Graph[VertexId, ED] = { require(maxSteps > 0, s"Maximum of steps must be greater than 0, but got ${maxSteps}") val lpaGraph = graph.mapVertices { case (vid, _) => vid } def sendMessage(e: EdgeTriplet[VertexId, ED]): Iterator[(VertexId, Map[VertexId, Long])] = { Iterator((e.srcId, Map(e.dstAttr -> 1L)), (e.dstId, Map(e.srcAttr -> 1L))) } def mergeMessage(count1: Map[VertexId, Long], count2: Map[VertexId, Long]) : Map[VertexId, Long] = { (count1.keySet ++ count2.keySet).map { i => val count1Val = count1.getOrElse(i, 0L) val count2Val = count2.getOrElse(i, 0L) i -> (count1Val + count2Val) }.toMap } def vertexProgram(vid: VertexId, attr: Long, message: Map[VertexId, Long]): VertexId = { if (message.isEmpty) attr else message.maxBy(_._2)._1 } val initialMessage = Map[VertexId, Long]() Pregel(lpaGraph, initialMessage, maxIterations = maxSteps)( vprog = vertexProgram, sendMsg = sendMessage, mergeMsg = mergeMessage) }
这里最终要的三个函数sendMessage、mergeMessage、vertexProgram,熟悉Pregel的同学可能就非常容易理解此处的意义了。
其中sendMessage向起点发送终点的标签,向终点发送起点的标签。这里说明了Graphx默认这里的图是无向的图的,如果有需要 的同学可以考虑改为有向图的情况下如何处理。
mergeMessage用于处理接收到的多个message的reduce,这里可以看到就是对同一个标签的合并,并且value值相加。
vertexProgram则是对node的操作,这里就是简单的把收到的message标签里面的值取出出现次数最多的nodeid用来更新本个node上的值。
