Skip to content

项目7:统计得分排名前10的网页——Spark GraphX 图计算框架 ​

先修基础:项目1-6(Spark概述、Scala基础、Spark Shell编程、Spark IDE编程、Spark SQL、Spark Streaming)


目录 ​


第一部分:项目背景与图计算概述 ​

1.1 项目背景 ​

为什么需要图计算? ​

互联网上海量网页互相链接,怎么判断一个网页的重要性? 社交网络中每个人有很多好友,怎么分析人际关系? 电商平台上用户和商品有购买关系,怎么做推荐?

这些问题都涉及到**"关系",而关系最好的表达方式就是图(Graph)**。

PageRank 就是最经典的图计算应用:

  • 谷歌搜索引擎的核心算法
  • 通过网页之间的链接关系计算网页的重要性(得分)
  • 得分高的网页排在搜索结果前面

项目场景 ​

某资讯网站有很多网页,网页之间互相有链接。 现有两份数据:

  • news_vertices.csv:网页信息(顶点数据)
  • news_edges.csv:网页之间的链接关系(边数据)

本项目要做的事:

  1. 用 GraphX 构建网页关系图
  2. 用 PageRank 算法计算每个网页的得分
  3. 找出得分排名前10的网页

1.2 什么是图计算? ​

图计算就是用**图(Graph)**这种数据结构来表示和计算"关系"。

常见应用场景 ​

场景顶点边
社交网络用户好友/关注关系
网页链接网页超链接
推荐系统用户、商品购买/浏览关系
交通网络城市/站点道路/航线
知识图谱实体(人、事、物)关系
金融风控账户转账关系

大厂的图计算 ​

  • 淘宝:图谱计算平台(商品推荐、用户画像)
  • 新浪微博:社交网络分析(粉丝关系、传播路径)
  • 腾讯:推荐应用(好友推荐、内容推荐)

💡 理解:

  • 以前学的 RDD、DataFrame 都是处理"独立的数据"
  • 图计算处理的是"数据之间的关系"
  • 有关系的数据,用图来算更高效、更直观

第二部分:图的基本概念 ​

2.1 图的定义 ​

图(Graph)通常表示为 G(V, E):

  • V(Vertex):顶点的集合,每个顶点代表一个对象
  • E(Edge):边的集合,每条边代表两个对象之间的关系

💡 形象比喻:

  • 顶点就像"人"
  • 边就像"人与人之间的关系"(朋友、同事、亲戚等)
  • 图就是"一群人和他们之间的关系网"

2.2 图的分类 ​

按边的方向分 ​

类型说明示例
有向图边有方向,A→B 不等于 B→A微博关注(A关注B,B不一定关注A)
无向图边没有方向,A-B 和 B-A 一样好友关系(A是B的朋友,B也是A的朋友)

按边的属性分 ​

类型说明示例
无权图边没有权重,只有"有/没有"关系是否是好友
有权图边有权重,表示关系的强度好友亲密度、距离

其他概念 ​

概念说明
完全图任意两个顶点之间都有边
简单图没有重复边,没有顶点到自身的边
度(Degree)一个顶点连接的边的数量
入度(In-degree)有向图中,指向该顶点的边数
出度(Out-degree)有向图中,从该顶点出发的边数

第三部分:Spark GraphX 基础概念 ​

3.1 什么是 GraphX? ​

GraphX 是 Spark 的图计算组件,是 Spark 四大组件之一。

  • 基于 RDD 实现的分布式图计算框架
  • 提供了丰富的图操作 API
  • 内置了常用的图算法(PageRank、最短路径等)

3.2 GraphX 中的核心概念 ​

顶点(Vertex) ​

  • 每个顶点有一个 VertexId(顶点ID,Long 类型)
  • 每个顶点有一个属性 VD(Vertex Attribute,可以是任意类型)
  • 表示:(VertexId, VD)

边(Edge) ​

  • 每条边有起点 srcId、目标点 dstId
  • 每条边有一个属性 ED(Edge Attribute)
  • Edge 类有三个字段:srcId、dstId、attr

图(Graph) ​

  • Graph 是 GraphX 的核心类
  • 包含顶点集合和边集合
  • 是图计算操作的入口

三种视图 ​

视图说明访问方式
顶点视图只看顶点graph.vertices
边视图只看边graph.edges
三元组视图顶点+边一起看graph.triplets

💡 三元组(Triplet):一条边 + 起点属性 + 目标点属性,是最完整的视图。


3.3 RDD 和 GraphX 的关系 ​

GraphX 概念底层 RDD说明
VertexRDD[VD]RDD[(VertexId, VD)]顶点RDD,VertexRDD继承自RDD
EdgeRDD[ED]RDD[Edge[ED]]边RDD,EdgeRDD继承自RDD
Graph[VD, ED]-图对象,包含顶点和边

💡 理解:GraphX 就是在 RDD 基础上封装了一层图的抽象, 底层还是 RDD,所以也有惰性求值、分区、持久化等特性。


第四部分:图的创建(3种方法) ​

4.1 准备工作:导入包 ​

scala
import org.apache.spark._
import org.apache.spark.graphx._
import org.apache.spark.rdd.RDD

4.2 方法一:Graph() —— 根据顶点和边创建 ​

最常用的方法,顶点和边都有属性。

语法 ​

scala
Graph(vertices: RDD[(VertexId, VD)], 
      edges: RDD[Edge[ED]], 
      defaultVertexAttr: VD = defaultValue)

参数说明 ​

参数类型说明
verticesRDD[(VertexId, VD)]顶点RDD,(顶点ID, 顶点属性)
edgesRDD[Edge[ED]]边RDD,Edge(起点ID, 目标点ID, 边属性)
defaultVertexAttrVD默认顶点属性(可选),边中出现但顶点数据中没有的顶点用这个值

示例:构建用户关系网络图 ​

顶点数据(vertices.txt):用户ID、姓名、职业

1 张三 工程师
2 李四 设计师
3 王五 产品经理
5 赵六 postdoc

边数据(edges.txt):起点ID、目标点ID、关系

1 2 朋友
2 3 同事
5 3 Advisor

代码:

scala
// 1. 读取顶点数据
val users = sc.textFile("/tipdm/data/vertices.txt").map { line =>
  val lines = line.split(" ")
  (lines(0).toLong, (lines(1), lines(2)))  // (VertexId, (姓名, 职业))
}

// 2. 读取边数据
val relationships = sc.textFile("/tipdm/data/edges.txt").map { line =>
  val lines = line.split(" ")
  Edge(lines(0).toLong, lines(1).toLong, lines(2))  // Edge(起点, 目标点, 边属性)
}

// 3. 默认顶点属性(边中出现但顶点数据中没有的顶点)
val defaultUser = ("John Doe", "Missing")

// 4. 创建图
val graph_urelate = Graph(users, relationships, defaultUser)

// 5. 查看顶点和边
graph_urelate.vertices.collect().foreach(println(_))
graph_urelate.edges.collect().foreach(println(_))

⚠️ 注意:VertexId 必须是 Long 类型,不是 Int! 这是新手最容易犯的错误之一。


4.3 方法二:Graph.fromEdges() —— 根据边创建 ​

只有边数据,顶点从边中自动提取,顶点属性是默认值。

语法 ​

scala
Graph.fromEdges(edges: RDD[Edge[ED]], defaultValue: VD)

示例 ​

scala
// 读取边数据
val records = sc.textFile("/tipdm/data/edges.txt")
val followers = records.map { x =>
  val fields = x.split(" ")
  Edge(fields(0).toLong, fields(1).toLong, fields(2))
}

// 根据边创建图,顶点属性默认值为 1L
val graph_fromEdges = Graph.fromEdges(followers, 1L)

// 查看
graph_fromEdges.vertices.collect().foreach(println(_))
graph_fromEdges.edges.collect().foreach(println(_))

💡 特点:

  • 只需要边数据,顶点自动从边中提取
  • 所有顶点的属性都是同一个默认值
  • 适合只关心结构、不关心顶点属性的场景

4.4 方法三:Graph.fromEdgeTuples() —— 根据边的二元组创建 ​

连边属性都不要,只有起点和目标点。

语法 ​

scala
Graph.fromEdgeTuples(rawEdges: RDD[(VertexId, VertexId)], 
                     defaultValue: VD,
                     uniqueEdges: Option[PartitionStrategy] = None)

示例 ​

scala
// 读取边数据,只取起点和目标点
val file = sc.textFile("/tipdm/data/edges.txt")
val edgesRDD = file.map(line => line.split(" "))
  .map(line => (line(0).toLong, line(1).toLong))  // (起点ID, 目标点ID)

// 创建图
val graph_fromEdgeTuples = Graph.fromEdgeTuples(edgesRDD, 1L)

// 查看
graph_fromEdgeTuples.vertices.collect().foreach(println(_))
graph_fromEdgeTuples.edges.collect().foreach(println(_))

💡 特点:

  • 最简单,只有顶点ID和边的连接关系
  • 边没有属性,顶点属性是默认值
  • 适合只关心拓扑结构的场景

4.5 三种创建方法对比 ​

方法顶点属性边属性适用场景
Graph(vertices, edges)✅ 自定义✅ 自定义最常用,顶点和边都有属性
Graph.fromEdges(edges, default)❌ 默认值✅ 自定义有边属性,顶点不重要
Graph.fromEdgeTuples(edges, default)❌ 默认值❌ 无只关心结构,最简单

第五部分:图的缓存与释放 ​

5.1 缓存 ​

和 RDD 一样,图也可以缓存,提高重复使用的效率。

方法 ​

scala
import org.apache.spark.storage.StorageLevel

// 方法1:cache(),默认内存存储
graph_urelate.cache()

// 方法2:persist(),默认内存存储
graph_fromEdges.persist()

// 方法3:persist(存储级别),指定存储级别
graph_fromEdgeTuples.persist(StorageLevel.MEMORY_ONLY)

常用存储级别 ​

级别说明
MEMORY_ONLY只存在内存
MEMORY_AND_DISK内存不够放磁盘
DISK_ONLY只放磁盘

5.2 释放缓存 ​

三种释放方式 ​

方法说明
graph.unpersist(blocking = false)释放整个图的缓存
graph.unpersistVertices(blocking = false)只释放顶点的缓存
graph.edges.unpersist(blocking = false)只释放边的缓存

blocking 参数 ​

  • blocking = false(默认):调用后立即返回,释放在后台执行
  • blocking = true:等待释放完成后再返回

示例 ​

scala
// 释放整个图
graph_fromEdges.unpersist(blocking = true)

// 只释放顶点
graph_fromEdges.unpersistVertices(blocking = true)

// 只释放边
graph_fromEdges.edges.unpersist(blocking = true)

⚠️ 注意:没有 graph.unpersistEdges() 这种方法! 释放边缓存要写成 graph.edges.unpersist()。


第六部分:图的数据查询 ​

6.1 基本信息查询 ​

顶点数和边数 ​

scala
// 顶点数量
graph_urelate.numVertices

// 边数量
graph_urelate.numEdges

6.2 三种视图 ​

① 顶点视图(vertices) ​

返回 VertexRDD[VD],即 RDD[(VertexId, VD)]。

scala
// 直接查看所有顶点
graph_urelate.vertices.collect().foreach(println(_))

// 用case匹配解构元组
graph_urelate.vertices.map {
  case (id, (name, prop)) => (prop, name)
}.collect().foreach(println(_))

// 过滤:职业是postdoc的顶点
graph_urelate.vertices.filter {
  case (id, (name, pos)) => pos == "postdoc"
}.collect().foreach(println(_))

// 用下标访问
graph_urelate.vertices.map { v =>
  (v._1, v._2._1, v._2._2)  // (ID, 姓名, 职业)
}.collect().foreach(println(_))

② 边视图(edges) ​

返回 EdgeRDD[ED],即 RDD[Edge[ED]]。 Edge 有三个字段:srcId(起点)、dstId(目标点)、attr(边属性)。

scala
// 直接查看所有边
graph_urelate.edges.collect().foreach(println(_))

// 过滤:起点ID大于目标点ID的边
graph_urelate.edges.filter {
  case Edge(src, dst, prop) => src > dst
}.collect().foreach(println(_))

// 用下标访问
graph_urelate.edges.map { e =>
  (e.attr, e.srcId, e.dstId)  // (边属性, 起点ID, 目标点ID)
}.collect().foreach(println(_))

③ 三元组视图(triplets) ​

最完整的视图,包含边 + 起点属性 + 目标点属性。 返回 RDD[EdgeTriplet[VD, ED]]。

EdgeTriplet 继承自 Edge,额外有:

  • srcAttr:起点的属性
  • dstAttr:目标点的属性
scala
// 直接查看所有三元组
graph_urelate.triplets.collect().foreach(println(_))

// 用下标访问
val graph_triplets = graph_urelate.triplets.map { triplet =>
  (
    triplet.srcAttr._1,  // 起点的姓名
    triplet.dstAttr._2,  // 目标点的职业
    triplet.srcId,       // 起点ID
    triplet.dstId        // 目标点ID
  )
}
graph_triplets.collect().foreach(println(_))

💡 理解:三元组就像"一条边的完整信息", 不仅知道谁连谁,还知道两边的顶点是什么样的。


6.3 度计算 ​

三种度 ​

方法说明返回类型
degrees总度数(入度+出度)VertexRDD[Int]
inDegrees入度数(指向该顶点的边数)VertexRDD[Int]
outDegrees出度数(从该顶点出发的边数)VertexRDD[Int]

示例 ​

scala
// 总度数
graph_urelate.degrees.collect()

// 入度数
graph_urelate.inDegrees.collect()

// 出度数
graph_urelate.outDegrees.collect()

度的统计:最大值、最小值、排序 ​

因为 degrees 返回的是 (VertexId, Int),要按度数排序需要先交换位置。

scala
// 交换位置,变成 (度数, 顶点ID)
val degree2 = graph_urelate.degrees.map(a => (a._2, a._1))

// 最大度数
print("max degree = " + (degree2.max()._2, degree2.max()._1))

// 最小度数
print("min degree = " + (degree2.min()._2, degree2.min()._1))

// 降序排序,取前3
degree2.sortByKey(true, 1).top(3).foreach(x => print(x._2, x._1))

💡 为什么要交换位置? 因为 sortByKey、top、max、min 这些方法都是对 Key(第一个元素) 操作的, 所以要把度数放到第一个位置。


第七部分:图的数据转换 ​

数据转换就是修改顶点或边的属性,不改变图的结构。

7.1 mapVertices() —— 修改顶点属性 ​

作用 ​

更新顶点的属性值,可以改变属性类型。

示例 ​

scala
// 原来顶点属性是 (姓名, 职业),现在只保留姓名
val newGraph = graph_urelate.mapVertices((id, attr) => attr._1)

// 查看新的顶点属性
newGraph.vertices.collect().foreach(println(_))

语法 ​

scala
graph.mapVertices((id, attr) => 新属性)
  • id:顶点ID
  • attr:原来的顶点属性
  • 返回值:新的顶点属性

7.2 mapEdges() —— 修改边属性 ​

作用 ​

更新边的属性值,可以改变属性类型。

示例 ​

scala
// 把边属性改成 起点ID+目标点ID
graph_urelate.mapEdges(e => e.srcId + e.dstId)
  .edges.collect().foreach(println(_))

语法 ​

scala
graph.mapEdges(e => 新属性)
  • e:Edge 对象(有 srcId、dstId、attr)
  • 返回值:新的边属性

7.3 mapTriplets() —— 用三元组修改边属性 ​

作用 ​

和 mapEdges 类似,但可以使用顶点的属性来修改边属性。

示例 ​

scala
// 把边属性替换成起点的姓名
val newGraph = graph_urelate.mapTriplets(triplet => triplet.srcAttr._1)

newGraph.triplets.collect().foreach(println(_))

语法 ​

scala
graph.mapTriplets(triplet => 新属性)
  • triplet:EdgeTriplet 对象(有 srcAttr、dstAttr、attr 等)
  • 返回值:新的边属性

7.4 三种 map 方法对比 ​

方法作用对象能访问什么修改什么
mapVertices顶点顶点ID + 顶点属性顶点属性
mapEdges边边的 srcId + dstId + attr边属性
mapTriplets边边 + 起点属性 + 目标点属性边属性

💡 记忆技巧:

  • 改顶点用 mapVertices
  • 改边用 mapEdges(只能看边本身)
  • 改边还要看顶点属性用 mapTriplets

第八部分:图的结构转换 ​

结构转换就是改变图的结构(顶点数量、边数量、边方向等)。

8.1 准备数据:用户社交网络图 ​

顶点(user.txt):用户ID、姓名、年龄

1,张三,25
2,李四,30
3,王五,35
4,赵六,28

边(relate.txt):起点ID、目标点ID、权重(粉丝支持度)

1,2,3
2,3,5
3,2,4
3,2,6
4,3,2

创建图:

scala
// 读取顶点
val users = sc.textFile("/tipdm/data/user.txt").map { line =>
  val lines = line.split(",")
  (lines(0).toLong, (lines(1), lines(2).toInt))
}

// 读取边
val relationships = sc.textFile("/tipdm/data/relate.txt").map { line =>
  val lines = line.split(",")
  Edge(lines(0).toLong, lines(1).toLong, lines(2).toInt)
}

// 创建图
val graph = Graph(users, relationships)
graph.triplets.collect().foreach(println(_))

8.2 reverse —— 反转边的方向 ​

作用 ​

返回一个新图,所有边的方向都反转了。

  • 顶点不变
  • 边的数量不变
  • 边的方向反过来

示例 ​

scala
val graph2 = graph.reverse

// 对比
graph.edges.collect()
graph2.edges.collect()

💡 理解:原来 A→B,反转后变成 B→A。 比如原来"张三关注李四",反转后变成"李四被张三关注"。


8.3 subgraph() —— 创建子图 ​

作用 ​

根据条件过滤顶点和边,生成子图。

语法 ​

scala
graph.subgraph(
  epred = edgePredicate,   // 边的过滤条件
  vpred = vertexPredicate  // 顶点的过滤条件
)

参数 ​

参数说明
epred边的谓词(过滤条件),返回 true 的边才保留
vpred顶点的谓词(过滤条件),返回 true 的顶点才保留

示例:边和顶点同时过滤 ​

scala
// 筛选:边权重大于3,且用户年龄大于30
val subGraph3 = graph.subgraph(
  epred = edge => edge.attr > 3,          // 边权重>3
  vpred = (id, attr) => attr._2 > 30     // 年龄>30
)

subGraph3.vertices.collect().foreach(println(_))
subGraph3.edges.collect().foreach(println(_))

💡 注意:两个参数都是可选的,可以只过滤边,也可以只过滤顶点。


8.4 mask() —— 取两个图的交集 ​

作用 ​

合并两个图,只保留两个图中都有的顶点和边。 结果图的属性来自调用 mask 的那个图(前一个图)。

示例 ​

scala
// graph 和 subGraph3 取交集,保留公共的顶点和边
val mask_graph = graph.mask(subGraph3)

mask_graph.triplets.collect().foreach(println(_))

💡 理解:mask 就像"两个图的交集", 只有两个图里都有的顶点和边才会保留。 属性用调用者(graph)的属性。


8.5 groupEdges() —— 合并重复边 ​

作用 ​

把同一起点到同一目标点的多条边合并成一条,边属性按指定方式合并。

前提 ​

必须先调用 partitionBy() 对图进行分区,确保相同的边在同一个分区。

常用分区策略 ​

PartitionStrategy.RandomVertexCut:随机顶点切分,把相同的边分到同一分区

示例 ​

scala
import org.apache.spark.graphx.PartitionStrategy

// 先分区
val graph2 = graph.partitionBy(PartitionStrategy.RandomVertexCut)

// 合并相同边,属性值相加
graph2.groupEdges((a, b) => a + b)
  .edges.collect().foreach(println(_))

💡 理解:如果有两条 3→2 的边,权重分别是4和6, groupEdges 后就变成一条 3→2 的边,权重是10。


8.6 结构转换方法总结 ​

方法作用改变顶点数改变边数
reverse反转边方向❌❌
subgraph过滤生成子图✅ 可能减少✅ 可能减少
mask两个图取交集✅ 可能减少✅ 可能减少
groupEdges合并重复边❌✅ 减少

第九部分:图的关联聚合操作(重点+难点) ​

关联聚合就是沿着边收集信息,然后聚合。 这是图计算最核心的操作,也是最难的部分。

9.1 collectNeighbors() —— 收集邻居顶点 ​

作用 ​

收集每个顶点的邻居顶点(ID + 属性)。

参数:方向 ​

参数说明
EdgeDirection.Out出边的邻居(以该顶点为起点的边的目标点)
EdgeDirection.In入边的邻居(以该顶点为目标点的边的起点)
EdgeDirection.Either所有邻居(入+出)

示例 ​

scala
import org.apache.spark.graphx.EdgeDirection

// 出边邻居
graph.collectNeighbors(EdgeDirection.Out).collect()

// 入边邻居
graph.collectNeighbors(EdgeDirection.In).collect()

// 所有邻居
graph.collectNeighbors(EdgeDirection.Either).collect()

⚠️ 注意:collectNeighbors 会把所有邻居数据拉到内存, 顶点度数很高时可能会 OOM,大数据量慎用。


9.2 collectNeighborIds() —— 收集邻居顶点ID ​

作用 ​

和 collectNeighbors 类似,但只返回顶点ID,不返回属性。 更轻量,占用内存更少。

示例 ​

scala
// 只收集邻居ID
graph.collectNeighborIds(EdgeDirection.Out).collect()

9.3 aggregateMessages() —— 核心聚合操作(重点) ​

作用 ​

GraphX 中最核心的聚合操作。

  • 每个边三元组可以向起点或目标点发送消息
  • 每个顶点把收到的所有消息聚合起来
  • 返回每个顶点的聚合结果

语法 ​

scala
graph.aggregateMessages[MsgType](
  sendMsg = triplet => { ... },   // 发送消息阶段(类似Map)
  mergeMsg = (a, b) => { ... },   // 合并消息阶段(类似Reduce)
  tripletFields = TripletFields.All  // 可选,指定需要哪些字段
)

三个参数 ​

参数说明
sendMsg对每个边三元组执行,发送消息给起点或目标点
mergeMsg对同一个顶点收到的多条消息进行合并
tripletFields指定需要访问哪些字段(优化性能)

sendMsg 中可用的方法 ​

方法说明
triplet.sendToSrc(msg)给起点发消息
triplet.sendToDst(msg)给目标点发消息

TripletFields 可选值 ​

值说明
TripletFields.Src只需要起点属性
TripletFields.Dst只需要目标点属性
TripletFields.All起点+目标点都需要(默认)

💡 指定 tripletFields 可以优化性能, 不需要的字段就不传输,减少网络开销。


示例:计算每个用户的追随者平均年龄 ​

需求:计算每个用户的粉丝(追随者)的平均年龄。

  • 追随者 = 边的起点(起点关注目标点)
  • 要算的是:每个目标点,它的所有起点的平均年龄
scala
// 1. aggregateMessages 聚合
val olderFollowers = graph.aggregateMessages[(Int, Int)](
  // sendMsg:给每个目标点发送 (1, 起点年龄)
  triplet => {
    triplet.sendToDst((1, triplet.srcAttr._2))
  },
  // mergeMsg:合并消息(人数相加,年龄总和相加)
  (a, b) => (a._1 + b._1, a._2 + b._2),
  TripletFields.All
)

// 2. 计算平均年龄
val avgAgeOfOlderFollowers: VertexRDD[Double] =
  olderFollowers.mapValues((id, value) => value match {
    case (count, totalAge) => totalAge / count.toDouble
  })

// 3. 查看结果
avgAgeOfOlderFollowers.collect().foreach(println(_))

过程详解 ​

sendMsg 阶段(每条边):
  边 1→2:给目标点2发送 (1, 25)  (1个粉丝,年龄25)
  边 2→3:给目标点3发送 (1, 30)  (1个粉丝,年龄30)
  边 3→2:给目标点2发送 (1, 35)  (1个粉丝,年龄35)
  ...

mergeMsg 阶段(每个顶点):
  顶点2收到:(1, 25) + (1, 35) = (2, 60)
  顶点3收到:(1, 30) + ... = ...

计算平均:
  顶点2:60 / 2 = 30.0  (粉丝平均年龄30岁)

💡 理解 aggregateMessages:

  • 就像"每个人给邻居递小纸条"
  • sendMsg 是"递纸条"的过程
  • mergeMsg 是"收到纸条后汇总"的过程
  • 最后每个人手里有一份汇总结果

9.4 joinVertices() —— 连接更新顶点属性 ​

作用 ​

把外部的顶点 RDD 和图的顶点连接,更新图的顶点属性。

  • 不能改变顶点属性的类型和个数
  • 外部 RDD 中没有的顶点,保持原值不变
  • 如果外部 RDD 有重复顶点,只有最后一个生效

语法 ​

scala
graph.joinVertices(verticesRDD) {
  case (id, oldAttr, newAttr) => 新属性  // 类型必须和原来一样
}

示例 ​

scala
// 1. 先用aggregateMessages计算每个顶点的粉丝数
val energys = graph.aggregateMessages[Int](
  triplet => triplet.sendToDst(1),  // 给目标点发1
  (a, b) => a + b                   // 求和
)

// 2. joinVertices:把粉丝数加到年龄上
val energys_name = graph.joinVertices(energys) {
  case (id, (name, age), energy) => (name, age + energy)
}

// 查看
energys.collect().foreach(println(_))
energys_name.vertices.collect().foreach(println(_))

⚠️ 注意:joinVertices 不能改变顶点属性的类型! 原来属性是 (String, Int),连接后也必须是 (String, Int)。


9.5 outerJoinVertices() —— 外连接更新顶点属性 ​

作用 ​

和 joinVertices 类似,但更灵活:

  • 可以改变顶点属性的类型和个数
  • 可以随意定义返回类型
  • 外部 RDD 中没有的顶点,用 None 处理

语法 ​

scala
graph.outerJoinVertices(verticesRDD) {
  case (id, oldAttr, Some(newAttr)) => 新属性  // 有匹配的情况
  case (id, oldAttr, None) => 新属性           // 没有匹配的情况
}

示例:把粉丝平均年龄加到顶点属性中 ​

scala
// 接上一节的 avgAgeOfOlderFollowers
val graph_avgAge = graph.outerJoinVertices(avgAgeOfOlderFollowers) {
  case (id, (name, age), Some(avgAge)) => (name, age, avgAge)  // 有粉丝的
  case (id, (name, age), None) => (name, age, 0.0)             // 没粉丝的
}

graph_avgAge.vertices.collect().foreach(println(_))

💡 对比:

  • joinVertices:内连接,属性类型不能变
  • outerJoinVertices:外连接,属性类型可以变,更灵活
  • 就像 SQL 的 inner join 和 left outer join 的区别

9.6 关联聚合方法总结 ​

方法作用返回类型
collectNeighbors收集邻居顶点(ID+属性)VertexRDD[Array[(VertexId, VD)]]
collectNeighborIds收集邻居顶点IDVertexRDD[Array[VertexId]]
aggregateMessages发消息+聚合(核心)VertexRDD[A]
joinVertices连接更新顶点(类型不变)Graph[VD, ED]
outerJoinVertices外连接更新顶点(类型可变)Graph[VD2, ED]

第十部分:项目实战——统计得分排名前10的网页 ​

10.0 数据说明 ​

数据文件 ​

  • news_vertices.csv:网页信息(顶点数据)
  • news_edges.csv:网页链接关系(边数据)

顶点字段(news_vertices.csv) ​

列索引字段说明
0网页ID顶点ID(Long)
1网页URL-
2网页标题顶点属性

边字段(news_edges.csv) ​

列索引字段说明
0源网页ID起点ID
1目标网页ID目标点ID

10.1 PageRank 算法简介 ​

什么是 PageRank? ​

PageRank(网页排名)是谷歌的核心算法,用来计算网页的重要性。

核心思想 ​

  • 一个网页被越多的其他网页链接到,它就越重要
  • 链接到它的网页越重要,它也越重要
  • 就像"投票",每个网页给它链接的网页投票,票的权重取决于自己的重要性

两种调用方式 ​

方式参数说明
静态pageRank(numIter)指定迭代次数,迭代完就停止
动态pageRank(tol)指定收敛阈值,结果变化小于阈值就停止

💡 动态调用更智能,收敛了就停,不用固定次数。


任务7.1:构建网页结构图 ​

需求 ​

读取网页数据和链接关系,构建图并缓存。

实现代码 ​

scala
// 1. 上传数据到HDFS
hdfs dfs -put news_vertices.csv /tipdm/data
hdfs dfs -put news_edges.csv /tipdm/data

// 2. 导入包
import org.apache.spark._
import org.apache.spark.graphx._
import org.apache.spark.rdd.RDD

// 3. 读取顶点数据(取第1列和第3列:ID和标题)
val vertice = sc.textFile("/tipdm/data/news_vertices.csv").map(
  line => {
    val data = line.split(",")
    (data(0).toLong, data(2))  // (网页ID, 网页标题)
  }
)

// 4. 读取边数据
val edge = sc.textFile("/tipdm/data/news_edges.csv").map(
  line => {
    val data = line.split(",")
    Edge(data(0).toLong, data(1).toLong, "连接")  // Edge(源网页, 目标网页, 边属性)
  }
)

// 5. 创建图并缓存
val graph = Graph(vertice, edge).cache()

// 6. 查看数据示例
graph.vertices.take(5).foreach(println(_))
graph.edges.take(5).foreach(println(_))

代码说明 ​

代码说明
data(0).toLong网页ID转成Long类型(VertexId必须是Long)
data(2)第3列(索引2)是网页标题,作为顶点属性
Edge(..., "连接")边属性设为"连接"(PageRank不关心边属性,随便设)
.cache()缓存图,后面PageRank要反复用

任务7.2:计算网页得分 ​

需求 ​

用 PageRank 算法计算每个网页的得分。

实现代码 ​

scala
// 动态调用:收敛阈值0.001
val web_rating = graph.pageRank(0.001)

// 查看得分
web_rating.vertices.take(10).foreach(println(_))

代码说明 ​

代码说明
graph.pageRank(0.001)动态PageRank,收敛阈值0.001
web_rating.vertices返回的图的顶点属性就是PageRank得分

💡 说明:

  • 参数 0.001 是收敛阈值,两次迭代结果变化小于这个值就停止
  • 返回值是一个新的 Graph,顶点属性变成了 PageRank 得分(Double类型)
  • 边属性也会变,变成权重

任务7.3:找出排名Top10的网页 ​

需求 ​

把网页标题和得分合并,按得分降序,取前10名。

实现代码 ​

scala
// 1. 取出PageRank的顶点(ID + 得分)
val web_score = web_rating.vertices

// 2. outerJoinVertices:把得分加到原图的顶点属性中
val graph_score = graph.outerJoinVertices(web_score) {
  case (id, title, Some(score)) => (title, score)  // 有得分的
  case (id, title, None) => (title, 0.0)           // 没得分的(理论上不会有)
}

// 3. 查看新图
graph_score.vertices.take(5).foreach(println(_))
graph_score.edges.take(5).foreach(println(_))

// 4. 按得分降序,取前10
graph_score.vertices.sortBy(_._2._2, false).take(10).foreach(println)

代码说明 ​

代码说明
outerJoinVertices(web_score)把得分和原图顶点连接
case (id, title, Some(score)) => (title, score)新属性是 (标题, 得分)
sortBy(_._2._2, false)按得分(_2._2)降序(false)排列
take(10)取前10名

💡 为什么用 outerJoinVertices? 因为 web_rating 的顶点属性只有得分,没有网页标题, 需要和原图的顶点属性(标题)合并,才能知道"哪个网页得多少分"。


任务7.4:完整程序实现(IDEA + spark-submit) ​

完整代码 ​

scala
import org.apache.spark.graphx._
import org.apache.spark.sql.SparkSession

object PageRank {
  def main(args: Array[String]): Unit = {
    // 1. 创建SparkSession
    val spark = SparkSession.builder()
      .appName("PageRank")
      .master("local[*]")
      .getOrCreate()
    val sc = spark.sparkContext
    
    // 2. 读取顶点数据
    val vertice = sc.textFile("/tipdm/data/news_vertices.csv")
      .map(line => {
        val data = line.split(",")
        (data(0).toLong, data(2))
      })
    
    // 3. 读取边数据
    val edge = sc.textFile("/tipdm/data/news_edges.csv")
      .map(line => {
        val data = line.split(",")
        Edge(data(0).toLong, data(1).toLong, "连接")
      })
    
    // 4. 创建图并缓存
    val graph = Graph(vertice, edge).cache()
    
    // 5. 计算PageRank
    val web_rating = graph.pageRank(0.001)
    
    // 6. 连接得分和顶点属性
    val web_score = web_rating.vertices
    val graph_score = graph.outerJoinVertices(web_score) {
      case (id, title, Some(score)) => (title, score)
      case (id, title, None) => (title, 0.0)
    }
    
    // 7. 取Top10
    val Top10 = graph_score.vertices.sortBy(_._2._2, false).take(10)
    
    // 8. 保存结果到HDFS
    sc.makeRDD(Top10).repartition(1).saveAsTextFile("/tipdm/data/Top10")
  }
}

打包提交 ​

bash
# 打包成page.jar后提交
spark-submit --master spark://master:7077 --class PageRank /opt/data/page.jar

查看结果 ​

bash
hdfs dfs -cat /tipdm/data/Top10/part-00000

第十一部分:常见问题与排错指南 ​

11.1 创建图类问题 ​

问题1:顶点ID类型不对报错 ​

现象:报类型不匹配错误 原因:VertexId 必须是 Long 类型,用了 Int 解决:顶点ID一定要 .toLong,不能是 .toInt

scala
// 错误写法
(data(0).toInt, data(2))

// 正确写法
(data(0).toLong, data(2))

问题2:边中出现的顶点在顶点数据中没有 ​

现象:报空指针或属性不对 原因:边数据中有些顶点ID在顶点数据中不存在 解决:创建图时指定 defaultVertexAttr(默认顶点属性)

scala
val graph = Graph(vertices, edges, defaultAttr)

11.2 操作类问题 ​

问题3:度排序不对 ​

现象:sortByKey 排序结果不对 原因:degrees 返回的是 (VertexId, Int),sortByKey 按ID排,不是按度数排 解决:先交换位置,把度数放到第一个位置

scala
// 错误:按ID排序
graph.degrees.sortByKey()

// 正确:按度数排序
graph.degrees.map(a => (a._2, a._1)).sortByKey()

问题4:groupEdges 不生效 ​

现象:调用 groupEdges 后重复边还在 原因:没有先调用 partitionBy,相同的边不在同一个分区 解决:先分区再合并

scala
import org.apache.spark.graphx.PartitionStrategy

// 错误
graph.groupEdges((a, b) => a + b)

// 正确
graph.partitionBy(PartitionStrategy.RandomVertexCut)
     .groupEdges((a, b) => a + b)

问题5:joinVertices 类型报错 ​

现象:报类型不匹配 原因:joinVertices 不能改变顶点属性的类型 解决:

  • 如果不需要改变类型,用 joinVertices
  • 如果需要改变类型(增加字段等),用 outerJoinVertices

11.3 性能类问题 ​

问题6:collectNeighbors 内存溢出 ​

现象:OOM(内存不足) 原因:顶点度数很高,收集的邻居数据太多 解决:

  • 尽量用 aggregateMessages 代替 collectNeighbors
  • 如果一定要用,先过滤再收集
  • 增加 executor 内存

问题7:PageRank 运行太慢 ​

现象:PageRank 迭代很多次,运行时间长 原因:收敛阈值太小,或者图太大 解决:

  • 适当调大收敛阈值(比如从0.001改成0.01)
  • 或者用静态调用,指定迭代次数
  • 确保图已经缓存

11.4 排错通用思路 ​

  1. 先看数据:vertices 和 edges 各 take 几行看看对不对
  2. 小数据测试:先用小数据验证逻辑
  3. 加 print 或 take:每一步都看看结果
  4. 注意类型:VertexId 是 Long,不是 Int
  5. 注意缓存:反复用的图要 cache()

第十二部分:实习 / 面试高频考点 ​

12.1 概念类(高频) ​

Q1:什么是图?什么是图计算? ​

图是由顶点(Vertex)和边(Edge)组成的数据结构,G(V, E)。 顶点代表对象,边代表对象之间的关系。 图计算就是用图来表示和计算"关系型数据"。 应用场景:社交网络、网页链接、推荐系统、交通网络、知识图谱等。

Q2:Spark GraphX 是什么? ​

GraphX 是 Spark 的图计算组件,基于 RDD 实现的分布式图计算框架。 提供了图的创建、查询、转换、聚合等操作, 还内置了 PageRank、最短路径等常用图算法。

Q3:GraphX 中的核心概念有哪些? ​

  • 顶点(Vertex):(VertexId, VD),VertexId 是 Long 类型
  • 边(Edge):srcId + dstId + attr
  • 图(Graph):包含顶点和边的集合
  • 三种视图:vertices(顶点视图)、edges(边视图)、triplets(三元组视图)
  • VertexRDD:顶点RDD,继承自 RDD[(VertexId, VD)]
  • EdgeRDD:边RDD,继承自 RDD[Edge[ED]]

Q4:创建图有哪几种方法?区别是什么? ​

三种方法:

  1. Graph(vertices, edges, defaultAttr):顶点和边都有属性,最常用
  2. Graph.fromEdges(edges, defaultValue):只有边,顶点属性是默认值
  3. Graph.fromEdgeTuples(edges, defaultValue):连边属性都没有,最简单 区别:顶点属性和边属性的有无、自定义程度不同。

Q5:图的三种视图是什么? ​

  1. 顶点视图(vertices):只看顶点,VertexRDD[VD]
  2. 边视图(edges):只看边,EdgeRDD[ED]
  3. 三元组视图(triplets):边 + 起点属性 + 目标点属性,最完整 EdgeTriplet 继承自 Edge,多了 srcAttr 和 dstAttr 两个属性。

Q6:什么是度?入度和出度的区别? ​

度(Degree)是一个顶点连接的边的数量。

  • 总度数(degrees):入度+出度
  • 入度(inDegrees):指向该顶点的边数
  • 出度(outDegrees):从该顶点出发的边数

Q7:mapVertices、mapEdges、mapTriplets 的区别? ​

  • mapVertices:修改顶点属性,只能访问顶点
  • mapEdges:修改边属性,只能访问边本身(srcId、dstId、attr)
  • mapTriplets:修改边属性,可以访问边 + 起点属性 + 目标点属性 前两个作用对象不同,mapTriplets 比 mapEdges 能访问更多信息。

Q8:subgraph 和 mask 的区别? ​

  • subgraph:根据过滤条件从一个图中筛选出子图
  • mask:取两个图的交集,保留两个图都有的顶点和边
  • subgraph 是"过滤",mask 是"两个图取交集"

Q9:joinVertices 和 outerJoinVertices 的区别? ​

  • joinVertices:内连接,不能改变顶点属性的类型和个数
  • outerJoinVertices:外连接,可以改变顶点属性的类型和个数,更灵活
  • 就像 SQL 的 inner join 和 left outer join 的区别。
  • 两者都是用来更新顶点属性的。

Q10:aggregateMessages 是什么?怎么用? ​

aggregateMessages 是 GraphX 核心的聚合操作。 两个阶段:

  1. sendMsg:每个边三元组向起点或目标点发送消息
  2. mergeMsg:每个顶点把收到的消息合并 返回 VertexRDD[A],每个顶点对应一个聚合结果。 可以指定 tripletFields 来优化性能(只传输需要的字段)。

12.2 原理类(中频) ​

Q11:GraphX 和 RDD 的关系? ​

GraphX 是基于 RDD 实现的。

  • VertexRDD 继承自 RDD[(VertexId, VD)]
  • EdgeRDD 继承自 RDD[Edge[ED]]
  • Graph 对象包含 VertexRDD 和 EdgeRDD 所以 GraphX 也有 RDD 的特性:惰性求值、分区、持久化等。

Q12:PageRank 算法的原理? ​

PageRank 是计算网页重要性的算法。 核心思想:

  • 一个网页被越多网页链接,越重要
  • 链接它的网页越重要,它也越重要 迭代计算,直到收敛(变化很小)。 GraphX 中可以直接调用 graph.pageRank(),有静态(指定次数)和动态(指定阈值)两种调用方式。

Q13:图的缓存和 RDD 的缓存一样吗? ​

类似,都是把数据存到内存中提高重复使用效率。 图的缓存有 cache() 和 persist() 两种。 释放可以释放整个图,也可以只释放顶点或只释放边。 注意:释放边是 graph.edges.unpersist(),不是 graph.unpersistEdges()。

Q14:groupEdges 为什么要先 partitionBy? ​

因为 groupEdges 需要把相同起点和目标点的边合并, 如果相同的边不在同一个分区,就没法合并。 partitionBy 按照分区策略(比如 RandomVertexCut)重新分区, 确保相同的边在同一个分区,才能正确合并。


12.3 实操类(高频) ​

Q15:VertexId 是什么类型? ​

Long 类型,不是 Int! 这是最常见的新手错误,创建顶点时一定要 .toLong。

Q16:怎么计算每个顶点的度数?怎么按度数排序? ​

度数:graph.degrees(总度数)、inDegrees(入度)、outDegrees(出度) 按度数排序:需要先把度数放到 key 的位置,因为 sortByKey 按 key 排。

scala
graph.degrees.map(a => (a._2, a._1)).sortByKey()

Q17:怎么用 PageRank? ​

scala
// 动态调用(推荐)
val ranks = graph.pageRank(0.001)

// 静态调用
val ranks = graph.pageRank(numIter = 10)

返回的是一个新图,顶点属性就是 PageRank 得分。

Q18:怎么把外部数据和图的顶点合并? ​

用 joinVertices 或 outerJoinVertices。

  • 不改变属性类型用 joinVertices
  • 要改变属性类型或增加字段用 outerJoinVertices 类似 SQL 的 join 操作。

附录:习题解析 ​

选择题解析 ​

1、答案:D 解析:图的顶点之间没有固定的连接数量限制,边数可以是任意的。

2、答案:C 解析:通过顶点数据和边数据创建图使用的是 Graph 方法,该方法需传入两个参数,第一个参数为顶点数据,第二个参数为边数据。

3、答案:A 解析:numVertices 和 numEdges 返回的分别是图中的顶点数量和边数量,vertices 和 edges 返回的分别是顶点 RDD 和边 RDD。

4、答案:B 解析:计算每个顶点的度数使用的是 degrees,首字母无需大写;inDegrees 和 outDegrees 分别用于计算每个顶点的入度数和出度数。

5、答案:D 解析:Graph.fromVertices() 方法不存在,且图必须包含边信息,仅凭顶点数据无法创建图。

6、答案:A 解析:joinVertices() 和 outerJoinVertices() 方法均可以用于更新顶点信息。 前者只能合并同类型,不能增加或删除属性,即属性个数和类型不可变;后者可以改变属性的类型和结构,属性类型可以改变。 aggregateMessages() 方法主要用于消息聚合,不直接更新顶点属性;collectNeighborsIds() 方法用于收集邻居顶点 ID,同样不更新属性。

7、答案:C 解析:mapVertices() 方法作用对象是顶点,仅用于转换顶点属性,不改变边;mapEdges() 方法作用对象是边,仅用于转换边属性,不改变顶点; mapTriplets() 方法作用对象是顶点和边,可以同时使用源顶点和目标顶点的属性来转换边属性;mapVerticesAndEdges() 方法不存在。

8、答案:C 解析:mask() 方法能够取两个图的公共顶点和边作新图,但保持的是前一个图顶点与边的属性。

9、答案:C 解析:Graph.unpersistEdges() 方法不存在,若想释放边缓存应该直接对边对象使用 unpersist() 方法,即 Graph.edges.unpersist()。

10、答案:D 解析:EdgeDirection.In 用于收集以邻居顶点作为起点、以该点作为目标点的邻居顶点的顶点 ID 和顶点信息; EdgeDirection.Out 用于收集以该顶点为起点、以邻居顶点为目标点的邻居顶点的顶点 ID 和顶点信息; EdgeDirection.Either 用于收集所有邻居顶点的顶点 ID 和顶点信息。


操作题解析 ​

题目 ​

英雄和武器的关系数据,查询"丘处机"使用的所有武器。

实现代码 ​

scala
import org.apache.spark.graphx.{Edge, Graph}
import org.apache.spark.rdd.RDD
import org.apache.spark.sql.SparkSession

object hero {
  def main(args: Array[String]): Unit = {
    // 1. 初始化SparkSession
    val spark = SparkSession.builder()
      .appName("hero")
      .master("local[*]")
      .getOrCreate()
    val sc = spark.sparkContext
    
    // 2. 读取英雄数据构建顶点RDD:(英雄ID, 英雄名称)
    val vertexRDD: RDD[(Long, String)] = sc.textFile("E:\\data\\hero.txt")
      .map(line => {
        val data = line.split(" ")
        (data(0).toLong, data(1))
      })
    
    // 3. 读取英雄-武器关系数据,去重
    val input: RDD[(String, String)] = sc.textFile("E:\\data\\hero_weapon.txt")
      .map(line => {
        val data = line.split(" ")
        ((data(0), data(1)), 1)
      })
      .groupBy(_._1)
      .map(_._1)
    
    // 4. 构建边RDD并创建图
    val edgeRDD: RDD[Edge[String]] = input.map(line => {
      Edge(line._1.toLong, line._2.toLong, "武器")
    })
    val graph = Graph(vertexRDD, edgeRDD)
    
    // 5. 查询英雄与其武器的对应关系,聚合,筛选"丘处机"
    graph.triplets
      .map(line => (line.srcAttr, line.dstAttr))
      .reduceByKey(_ + "," + _)
      .filter(_._1 == "丘处机")
      .collect
      .foreach(println(_))
  }
}

代码说明 ​

代码说明
vertexRDD顶点:英雄ID → 英雄名称
input英雄-武器关系,先 groupBy 去重
edgeRDD边:英雄ID → 武器ID,边属性"武器"
graph.triplets三元组视图,能同时看到起点(英雄)和目标点(武器)的属性
.map(line => (line.srcAttr, line.dstAttr))提取(英雄名,武器名)
.reduceByKey(_ + "," + _)按英雄聚合,多个武器用逗号拼接
.filter(_._1 == "丘处机")筛选丘处机

💡 思路:

  1. 英雄和武器分别是图的两类顶点
  2. 英雄和武器之间的"使用关系"是边
  3. 用 triplets 视图拿到英雄名和武器名
  4. 按英雄聚合(reduceByKey)
  5. 筛选出丘处机

笔记版本:V1.0 对应教材:《Spark大数据技术与应用(第3版)》人民邮电出版社 对应项目:项目7 统计得分排名前10的网页——Spark GraphX图计算框架 最后更新:2026年8月

基于 Vite 强力驱动 | 纯静态轻量托管