3.3.4 Spark GraphX框架及编程实例
Spark GraphX是一个分布式图处理框架,它是基于Spark平台对图计算和图挖掘提供的简洁易用组件,极大地方便了对分布式图处理的需求。
众所周知,社交网络中人与人之间有很多关系链,例如Twitter、Facebook、微博和微信等,这些都是大数据产生的地方,都需要图计算。同时,现在的图处理基本都是分布式的图处理,而并非单机处理。Spark GraphX由于底层是基于Spark来处理的,所以天然就是一个分布式的图处理系统。
GraphX通过扩展Spark RDD引入了一个新的图抽象数据结构,一个将有效信息放入顶点和边的有向多重图。如同Spark的每一个模块一样,这些模块都有一个基于RDD的便于自己计算的抽象数据结构(如SQL的DataFrame,Streaming的DStream)。为了方便图计算,GraphX公开了一系列基本运算(InDegress,OutDegress,subgraph等),也有许多用于图计算算法工具包(PageRank,TriangleCount,ConnectedComponents等算法)。相对于其他分布式图计算框架,GraphX最大的贡献,也是大多数开发者喜欢它的原因,是它在Spark之上提供了一站式解决方案,可以方便且高效地完成图计算的一整套流水作业,即在实际开发中,可以使用核心模块来完成海量数据的清洗与分析阶段,用SQL模块来打通与数据仓库的通道,用Streaming打造实时流处理通道,用GraphX图计算算法来对网页中复杂的业务关系进行计算,最后使用MLlib以及SparkR来完成数据挖掘算法处理。
Spark的每一个模块,都有自己的抽象数据结构,GraphX的核心抽象是弹性分布式属性图(Resilient Distribute Property Graph),一种点和边都带有属性的有向多重图。GraphX同时拥有Table和Graph两种视图,但只需一种物理存储,这两种视图都有自己独有的操作符,从而获得灵活的操作和较高的执行效率,GraphX结构图如图3⁃22所示。

图3⁃22 GraphX结构图
对Graph视图的所有操作,最终都会转换成其关联的Table视图的RDD操作来完成。这样对一个图的计算,最终在逻辑上,等价于一系列RDD的转换过程。因此,Graph最终具备了RDD的3个关键特性:不可变的、分布式的和容错的,其中最关键的是不可变的。逻辑上,所有图的转换和操作都产生了一个新图。物理上,GraphX会有一定程度的不变顶点和边的复用优化,对用户透明。
两种视图底层共用的物理数据,由RDD[VertexPartition]和RDD[EdgePartition]这两个RDD组成。点和边实际都不是以表Collection[tuple]的形式存储的,而是由VertexPartition/EdgePartition在内部存储一个带索引结构的分片数据块,以加速不同视图下的遍历速度。不变的索引结构在RDD转换过程中是共用的,降低了计算和存储开销。
图的分布式存储采用点分割模式,而且使用partitionBy方法,由用户指定不同的划分策略(Partition Strategy)。划分策略会将边分配到各个EdgePartition,顶点Master分配到各个VertexPartition,EdgePartition也会缓存本地边关联点的Ghost副本。划分策略的不同会影响到所需要缓存的Ghost副本数量,以及每个EdgePartition分配的边的均衡程度,需要根据图的结构特征选取最佳策略。
图3⁃23为GraphX架构图,可以看出,GraphX的整体架构分为三个部分。

图3⁃23 GraphX架构
①存储层和原语层:Graph类是图计算的核心类,内部含有VertexRDD、EdgeRDD和RDD[EdgeTriplet]引用。GraphImpl是Graph类的子类,实现了图操作。
②接口层:在底层RDD的基础之上实现Pregel模型,是BSP模式的计算接口。
③算法层:基于Pregel接口实现了常用的图算法,包含:PageRank、SVDPlusPlus、TriangleCount、ConnectedComponents、StronglyConnectedConponents等算法。
GraphX内部提供了三种RDD来对一个有向多重图的属性进行描述,VertexRDD和EdgeRDD都继承于RDD。属性图的每一个顶点由具有64位长度的唯一标识符(VertexID)作为主键,GraphX并没有对顶点添加任何顺序的约束,每一条边都有相应的源顶点(srcVertex)和目标顶点(dstVertex)。对于属性图的改变是通过生成新的图来完成,原图的主要部分(属性和索引)会被重用,通过启发式执行顶点分区,在不同的执行器中进行顶点的划分。与RDD一样,当系统发生故障的时候,图中的每个分区都可以重建。
VertexRDD[A]表示具有A属性的顶点集合。在内部将顶点的属性集合使用一个可重复使用的暗哈希表来存储,因此当两个VertexRDD继承自同一个VertexRDD的时候,它们可以在常数时间内进行合并操作,而不需要重新去计算哈希值。
EdgeRDD[ED,VD]继承自RDD[Edge[ED]],以各种分区策略(PartitionStrategy)将边划分为不同的块。在每个分区中,边属性和邻接结构分别存储,这使得更改属性值的时候能够实现最大限度的复用。
除了上面两种RDD,还有一种RDD结构为Triplet,它相当于在EdgeRDD的结构上加上了点的属性。同VertexRDD和EdgeRDD类似的,它也提供了多种常用的操作。消息通过边triplet的一个函数被并行计算,消息的计算既会访问源顶点特征也会访问目的顶点特征。(https://www.daowen.com)
可见,相比较Pregel,GraphX具有以下优点:
①允许用户把数据当作一个图和一个集合(RDD),而不需要数据移动或者复制。
②Spark GraphX可以无缝与Spark SQL、MLlib等结合,方便且高效地完成图计算整套流水作业。
如同Spark一样,GraphX的Graph类提供了丰富的图运算符,大致结构如图3⁃24所示。可以在官方GraphX Programming Guide中找到每个函数的详细说明。

图3⁃24 GraphX编程接口
著名的PageRank算法就是基于GraphX框架实现的。PageRank算法即网页排名算法,是Google创始人拉里·佩奇和谢尔盖·布林于1997年构建早期的搜索系统原型时提出的链接分析算法。自从Google在商业上获得巨大成功后,该算法引起了研究者们广泛关注,其中很多重要的链接算法都是在PageRank算法基础上衍生出来的。PageRank算法是Google用来标识网页等级的重要依据,也是Google用来衡量一个站点的好坏的唯一标准。
对网页进行排名,需要有量化的依据,因此PageRank算法对每一个网页进行计算,得到一个在0到10范围内的值,即该网页的PageRank值,简称PR值。PR值越高说明该网页越受欢迎(越重要)。
PageRank算法的核心步骤如下:
第一步 初始化。
PageRank算法基于两个假设:数量假设和质量假设。首先通过链接关系构建Web图,如图3⁃25所示。网络中每个页面对应Web图中的一个顶点,若网页A中包含一条指向B的链接,则Web图中存在一条由顶点A指向顶点B的边。然后为Web图中的每个顶点设置初始的PR值(通常设定每个顶点的初始PR值为1/N,其中N为网络中网页的个数)。

图3⁃25 Web图
第二步 迭代计算。
假设一个用户在访问某网页时,将其跳转到该网页上各超链接页面的概率相同。如图3⁃25所示,网页A链向网页B、C、D,所以根据假设,用户从A跳转到B、C、D的概率各为1/3。PageRank算法在每一轮迭代计算的过程中,首先将每个网页当前的PR值平均分配到该网页指向的超链接页面上,这样每个网页便获得了相应的权值。然后将这些权值进行求和,最后得到该网页的新PR值。当每个网页的PR值都获得更新后,就完成了一轮迭代。
第三步 结束迭代。
随着每一轮的迭代计算,网页中PR值会不断得到更新。当迭代达到一定次数,或者每个网页的PR值固定不变,再或者PR值收敛至某一范围内时,该算法停止。算法停止时每个网页的PR值就是该网页最终的PR值。
GraphX提供了静态和动态PageRank的实现方法,这些方法在PageRank对象中。静态的PageRank运行固定次数的迭代,而动态的PageRank一直运行直到收敛为止。