Spark系列 (七)SparkGraphX下的Pregel方法

文章目录

  • Pregel框架:

  • 一:Spark GraphX Pregel:

  • 二:Pregel计算过程:

  • Pregel函数源码及各个参数解析:

  • 三:案例:单源最短路径

  • 第一步:调用pregel方法:

  • 第二步:第一次迭代:

  • 第三步:第二次迭代:

  • 第四步:不断迭代,直至所有顶点处于钝化态

  • 案例代码如下:

Pregel框架:

一:Spark GraphX Pregel:

  • Pregel是google提出的用于大规模分布式图计算框架
    • 图遍历(bfs)
    • 单源最短路径(sssp)
    • pageRank计算
  • Pregel的计算有一系列迭代组成
  • Pregel迭代过程
    • 每个顶点从上一个superstep接收入站消息
    • 计算顶点新的属性
    • 在下一个superstep中向相邻的顶点发送消息
    • 当没有剩余消息时,迭代结束

二:Pregel计算过程:

Pregel函数源码及各个参数解析:

1def pregel[A: ClassTag]( 2 // 图节点的初始信息 3 initialMsg: A, 4 // 最大迭代次数 5 maxIterations: Int = Int.MaxValue, 6 // 7 activeDirection: EdgeDirection = EdgeDirection.Either)( 8 vprog: (VertexId, VD, A) => VD, 9 sendMsg: EdgeTriplet[VD, ED] => Iterator[(VertexId, A)], 10 mergeMsg: (A, A) => A) 11 : Graph[VD, ED] = { 12 Pregel(graph, initialMsg, maxIterations, activeDirection)(vprog, sendMsg, mergeMsg) 13 } 14

参数

说明

initialMsg

图初始化的时候,开始模型计算的时候,所有节点都会先收到一个消息

maxIterations

最大迭代次数

activeDirection

规定了发送消息的方向

vprog

节点调用该消息将聚合后的数据和本节点进行属性的合并

sendMsg

激活态的节点调用该方法发送消息

mergeMsg

如果一个节点接收到多条消息,先用mergeMsg 来将多条消息聚合成为一条消息,如果节点只收到一条消息,则不调用该函数

三:案例:单源最短路径

首先要清楚关于 顶点 的两点知识:

  1. 顶点 的状态有两种:
    (1)、钝化态【类似于休眠,不做任何事】
    (2)、激活态【干活】

  2. 顶点 能够处于激活态需要有条件:
    (1)、成功收到消息 或者
    (2)、成功发送了任何一条消息

第一步:调用pregel方法:

从5出发,除自身顶点外所有顶点都将接收一条初始消息initialMsg,使所有顶点处于激活态,并将属性改成无穷大。自身顶点为0。

第二步:第一次迭代:

所有顶点以EdgeDirection.Out的边方向调用sendMsg方法发送消息给目标顶点,如果 源顶点的属性+边的属性<目标顶点的属性,则发送消息。否则不发送。

之后只有两条边的信息发送成功了

5—>3(0+8<Double.Infinity , 成功),
5—>6(0+3<Double.Infinity , 成功)

此时只有5,3,6处于激活态了,3,6调用vprog方法,将属性合并。

第三步:第二次迭代:

处于激活态的3,6调用sendMsg方法发送消息。

最后只有3—>2(8+4<Double.Infinity,成功)

此时只有3,2处于激活态,2调用vprog方法将属性合并。

第四步:不断迭代,直至所有顶点处于钝化态

每个顶点的属性,就是顶点5到达各个顶点的最短距离。

案例代码如下:

1package com.wyw 2 import org.apache.spark.{SparkConf, SparkContext} 3 import org.apache.spark.graphx._ 4 import org.apache.spark.rdd.RDD 5object Pregel { 6 7 //1、创建SparkContext 8 val sparkConf = new SparkConf().setAppName("GraphxHelloWorld").setMaster("local[*]") 9 val sparkContext = new SparkContext(sparkConf) 10 11 //2、创建顶点 12 val vertexArray = Array( 13 (1L, ("Alice", 28)), 14 (2L, ("Bob", 27)), 15 (3L, ("Charlie", 65)), 16 (4L, ("David", 42)), 17 (5L, ("Ed", 55)), 18 (6L, ("Fran", 50)) 19 ) 20 val vertexRDD: RDD[(VertexId, (String,Int))] = sparkContext.makeRDD(vertexArray) 21 22 //3、创建边,边的属性代表 相邻两个顶点之间的距离 23 val edgeArray = Array( 24 Edge(2L, 1L, 7), 25 Edge(2L, 4L, 2), 26 Edge(3L, 2L, 4), 27 Edge(3L, 6L, 3), 28 Edge(4L, 1L, 1), 29 Edge(2L, 5L, 2), 30 Edge(5L, 3L, 8), 31 Edge(5L, 6L, 3) 32 ) 33 val edgeRDD: RDD[Edge[Int]] = sparkContext.makeRDD(edgeArray) 34 35 36 //4、创建图(使用aply方式创建) 37 val graph1 = Graph(vertexRDD, edgeRDD) 38 39 /* ************************** 使用pregle算法计算 ,顶点5 到 各个顶点的最短距离 ************************** */ 40 41 //被计算的图中 起始顶点id,初始化把点属性全部换成正无穷 42 val srcVertexId = 5L 43 val initialGraph = graph1.mapVertices{ 44 case (vid,(name,age)) => 45 if (vid==srcVertexId) 46 0.0 47 else 48 Double.PositiveInfinity 49 } 50 51 //5、调用pregel柯里化函数 52 val pregelGraph: Graph[Double, PartitionID] = initialGraph.pregel( 53 Double.PositiveInfinity, 54 Int.MaxValue, 55 EdgeDirection.Out 56 )( 57 // 传三个匿名函数参数 58 // 我收到消息后与本节点判断 59 (vid: VertexId, vd: Double, distMsg: Double) => { 60 // 比较两者值 61 val minDist = math.min(vd, distMsg) 62 println(s"顶点$vid,属性$vd,收到消息$distMsg,合并后的属性$minDist") 63 // 把小数据发送出去 64 minDist 65 }, 66 // 是不是要向下个点发数据 67 (edgeTriplet: EdgeTriplet[Double,PartitionID]) => { 68 // 检查起点+权重的值 和终点的值判断,小于才发送 69 if (edgeTriplet.srcAttr + edgeTriplet.attr < edgeTriplet.dstAttr) { 70 println(s"顶点${edgeTriplet.srcId} 给 顶点${edgeTriplet.dstId} 发送消息 ${edgeTriplet.srcAttr + edgeTriplet.attr}") 71 72 Iterator[(VertexId, Double)]((edgeTriplet.dstId, edgeTriplet.srcAttr + edgeTriplet.attr)) 73 } else { 74 Iterator.empty 75 } 76 }, 77 // 多个消息进行判断,取最小的消息发送,每次都处理2个 78 (msg1: Double, msg2: Double) => math.min(msg1, msg2) 79 ) 80 81 //6、输出结果 82 // pregelGraph.triplets.collect().foreach(println) 83 // println(pregelGraph.vertices.collect.mkString("\n")) 84 85 //7、关闭SparkContext 86 sparkContext.stop() 87 88}
点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

swap空间的增减方法

(1)增大swap空间去激活swap交换区:swapoff v /dev/vg00/lvswap扩展交换lv:lvextend L 10G /dev/vg00/lvswap重新生成swap交换区:mkswap /dev/vg00/lvswap激活新生成的交换区:swapon v /dev/vg00/lvswap

Java获得今日零时零分零秒的时间(Date型)

publicDatezeroTime()throwsParseException{    DatetimenewDate();    SimpleDateFormatsimpnewSimpleDateFormat("yyyyMMdd00:00:00");    SimpleDateFormatsimp2newS