项目7:统计得分排名前10的网页——Spark GraphX 图计算框架
先修基础:项目1-6(Spark概述、Scala基础、Spark Shell编程、Spark IDE编程、Spark SQL、Spark Streaming)
目录
- 第一部分:项目背景与图计算概述
- 第二部分:图的基本概念
- 第三部分:Spark GraphX 基础概念
- 第四部分:图的创建(3种方法)
- 第五部分:图的缓存与释放
- 第六部分:图的数据查询
- 第七部分:图的数据转换
- 第八部分:图的结构转换
- 第九部分:图的关联聚合操作(重点+难点)
- 第十部分:项目实战——统计得分排名前10的网页
- 第十一部分:常见问题与排错指南
- 第十二部分:实习 / 面试高频考点
- 附录:习题解析
第一部分:项目背景与图计算概述
1.1 项目背景
为什么需要图计算?
互联网上海量网页互相链接,怎么判断一个网页的重要性? 社交网络中每个人有很多好友,怎么分析人际关系? 电商平台上用户和商品有购买关系,怎么做推荐?
这些问题都涉及到**"关系",而关系最好的表达方式就是图(Graph)**。
PageRank 就是最经典的图计算应用:
- 谷歌搜索引擎的核心算法
- 通过网页之间的链接关系计算网页的重要性(得分)
- 得分高的网页排在搜索结果前面
项目场景
某资讯网站有很多网页,网页之间互相有链接。 现有两份数据:
news_vertices.csv:网页信息(顶点数据)news_edges.csv:网页之间的链接关系(边数据)
本项目要做的事:
- 用 GraphX 构建网页关系图
- 用 PageRank 算法计算每个网页的得分
- 找出得分排名前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 准备工作:导入包
import org.apache.spark._
import org.apache.spark.graphx._
import org.apache.spark.rdd.RDD4.2 方法一:Graph() —— 根据顶点和边创建
最常用的方法,顶点和边都有属性。
语法
Graph(vertices: RDD[(VertexId, VD)],
edges: RDD[Edge[ED]],
defaultVertexAttr: VD = defaultValue)参数说明
| 参数 | 类型 | 说明 |
|---|---|---|
vertices | RDD[(VertexId, VD)] | 顶点RDD,(顶点ID, 顶点属性) |
edges | RDD[Edge[ED]] | 边RDD,Edge(起点ID, 目标点ID, 边属性) |
defaultVertexAttr | VD | 默认顶点属性(可选),边中出现但顶点数据中没有的顶点用这个值 |
示例:构建用户关系网络图
顶点数据(vertices.txt):用户ID、姓名、职业
1 张三 工程师
2 李四 设计师
3 王五 产品经理
5 赵六 postdoc边数据(edges.txt):起点ID、目标点ID、关系
1 2 朋友
2 3 同事
5 3 Advisor代码:
// 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() —— 根据边创建
只有边数据,顶点从边中自动提取,顶点属性是默认值。
语法
Graph.fromEdges(edges: RDD[Edge[ED]], defaultValue: VD)示例
// 读取边数据
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() —— 根据边的二元组创建
连边属性都不要,只有起点和目标点。
语法
Graph.fromEdgeTuples(rawEdges: RDD[(VertexId, VertexId)],
defaultValue: VD,
uniqueEdges: Option[PartitionStrategy] = None)示例
// 读取边数据,只取起点和目标点
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 一样,图也可以缓存,提高重复使用的效率。
方法
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:等待释放完成后再返回
示例
// 释放整个图
graph_fromEdges.unpersist(blocking = true)
// 只释放顶点
graph_fromEdges.unpersistVertices(blocking = true)
// 只释放边
graph_fromEdges.edges.unpersist(blocking = true)⚠️ 注意:没有
graph.unpersistEdges()这种方法! 释放边缓存要写成graph.edges.unpersist()。
第六部分:图的数据查询
6.1 基本信息查询
顶点数和边数
// 顶点数量
graph_urelate.numVertices
// 边数量
graph_urelate.numEdges6.2 三种视图
① 顶点视图(vertices)
返回 VertexRDD[VD],即 RDD[(VertexId, VD)]。
// 直接查看所有顶点
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(边属性)。
// 直接查看所有边
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:目标点的属性
// 直接查看所有三元组
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] |
示例
// 总度数
graph_urelate.degrees.collect()
// 入度数
graph_urelate.inDegrees.collect()
// 出度数
graph_urelate.outDegrees.collect()度的统计:最大值、最小值、排序
因为 degrees 返回的是 (VertexId, Int),要按度数排序需要先交换位置。
// 交换位置,变成 (度数, 顶点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() —— 修改顶点属性
作用
更新顶点的属性值,可以改变属性类型。
示例
// 原来顶点属性是 (姓名, 职业),现在只保留姓名
val newGraph = graph_urelate.mapVertices((id, attr) => attr._1)
// 查看新的顶点属性
newGraph.vertices.collect().foreach(println(_))语法
graph.mapVertices((id, attr) => 新属性)id:顶点IDattr:原来的顶点属性- 返回值:新的顶点属性
7.2 mapEdges() —— 修改边属性
作用
更新边的属性值,可以改变属性类型。
示例
// 把边属性改成 起点ID+目标点ID
graph_urelate.mapEdges(e => e.srcId + e.dstId)
.edges.collect().foreach(println(_))语法
graph.mapEdges(e => 新属性)e:Edge 对象(有 srcId、dstId、attr)- 返回值:新的边属性
7.3 mapTriplets() —— 用三元组修改边属性
作用
和 mapEdges 类似,但可以使用顶点的属性来修改边属性。
示例
// 把边属性替换成起点的姓名
val newGraph = graph_urelate.mapTriplets(triplet => triplet.srcAttr._1)
newGraph.triplets.collect().foreach(println(_))语法
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创建图:
// 读取顶点
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 —— 反转边的方向
作用
返回一个新图,所有边的方向都反转了。
- 顶点不变
- 边的数量不变
- 边的方向反过来
示例
val graph2 = graph.reverse
// 对比
graph.edges.collect()
graph2.edges.collect()💡 理解:原来 A→B,反转后变成 B→A。 比如原来"张三关注李四",反转后变成"李四被张三关注"。
8.3 subgraph() —— 创建子图
作用
根据条件过滤顶点和边,生成子图。
语法
graph.subgraph(
epred = edgePredicate, // 边的过滤条件
vpred = vertexPredicate // 顶点的过滤条件
)参数
| 参数 | 说明 |
|---|---|
epred | 边的谓词(过滤条件),返回 true 的边才保留 |
vpred | 顶点的谓词(过滤条件),返回 true 的顶点才保留 |
示例:边和顶点同时过滤
// 筛选:边权重大于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 的那个图(前一个图)。
示例
// graph 和 subGraph3 取交集,保留公共的顶点和边
val mask_graph = graph.mask(subGraph3)
mask_graph.triplets.collect().foreach(println(_))💡 理解:mask 就像"两个图的交集", 只有两个图里都有的顶点和边才会保留。 属性用调用者(graph)的属性。
8.5 groupEdges() —— 合并重复边
作用
把同一起点到同一目标点的多条边合并成一条,边属性按指定方式合并。
前提
必须先调用 partitionBy() 对图进行分区,确保相同的边在同一个分区。
常用分区策略
PartitionStrategy.RandomVertexCut:随机顶点切分,把相同的边分到同一分区
示例
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 | 所有邻居(入+出) |
示例
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,不返回属性。 更轻量,占用内存更少。
示例
// 只收集邻居ID
graph.collectNeighborIds(EdgeDirection.Out).collect()9.3 aggregateMessages() —— 核心聚合操作(重点)
作用
GraphX 中最核心的聚合操作。
- 每个边三元组可以向起点或目标点发送消息
- 每个顶点把收到的所有消息聚合起来
- 返回每个顶点的聚合结果
语法
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 可以优化性能, 不需要的字段就不传输,减少网络开销。
示例:计算每个用户的追随者平均年龄
需求:计算每个用户的粉丝(追随者)的平均年龄。
- 追随者 = 边的起点(起点关注目标点)
- 要算的是:每个目标点,它的所有起点的平均年龄
// 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 有重复顶点,只有最后一个生效
语法
graph.joinVertices(verticesRDD) {
case (id, oldAttr, newAttr) => 新属性 // 类型必须和原来一样
}示例
// 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 处理
语法
graph.outerJoinVertices(verticesRDD) {
case (id, oldAttr, Some(newAttr)) => 新属性 // 有匹配的情况
case (id, oldAttr, None) => 新属性 // 没有匹配的情况
}示例:把粉丝平均年龄加到顶点属性中
// 接上一节的 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 | 收集邻居顶点ID | VertexRDD[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:构建网页结构图
需求
读取网页数据和链接关系,构建图并缓存。
实现代码
// 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 算法计算每个网页的得分。
实现代码
// 动态调用:收敛阈值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名。
实现代码
// 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)
完整代码
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")
}
}打包提交
# 打包成page.jar后提交
spark-submit --master spark://master:7077 --class PageRank /opt/data/page.jar查看结果
hdfs dfs -cat /tipdm/data/Top10/part-00000第十一部分:常见问题与排错指南
11.1 创建图类问题
问题1:顶点ID类型不对报错
现象:报类型不匹配错误 原因:VertexId 必须是 Long 类型,用了 Int 解决:顶点ID一定要 .toLong,不能是 .toInt
// 错误写法
(data(0).toInt, data(2))
// 正确写法
(data(0).toLong, data(2))问题2:边中出现的顶点在顶点数据中没有
现象:报空指针或属性不对 原因:边数据中有些顶点ID在顶点数据中不存在 解决:创建图时指定 defaultVertexAttr(默认顶点属性)
val graph = Graph(vertices, edges, defaultAttr)11.2 操作类问题
问题3:度排序不对
现象:sortByKey 排序结果不对 原因:degrees 返回的是 (VertexId, Int),sortByKey 按ID排,不是按度数排 解决:先交换位置,把度数放到第一个位置
// 错误:按ID排序
graph.degrees.sortByKey()
// 正确:按度数排序
graph.degrees.map(a => (a._2, a._1)).sortByKey()问题4:groupEdges 不生效
现象:调用 groupEdges 后重复边还在 原因:没有先调用 partitionBy,相同的边不在同一个分区 解决:先分区再合并
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 排错通用思路
- 先看数据:vertices 和 edges 各 take 几行看看对不对
- 小数据测试:先用小数据验证逻辑
- 加 print 或 take:每一步都看看结果
- 注意类型:VertexId 是 Long,不是 Int
- 注意缓存:反复用的图要 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:创建图有哪几种方法?区别是什么?
三种方法:
Graph(vertices, edges, defaultAttr):顶点和边都有属性,最常用Graph.fromEdges(edges, defaultValue):只有边,顶点属性是默认值Graph.fromEdgeTuples(edges, defaultValue):连边属性都没有,最简单 区别:顶点属性和边属性的有无、自定义程度不同。
Q5:图的三种视图是什么?
- 顶点视图(vertices):只看顶点,VertexRDD[VD]
- 边视图(edges):只看边,EdgeRDD[ED]
- 三元组视图(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 核心的聚合操作。 两个阶段:
- sendMsg:每个边三元组向起点或目标点发送消息
- 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 排。
scalagraph.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 和顶点信息。
操作题解析
题目
英雄和武器的关系数据,查询"丘处机"使用的所有武器。
实现代码
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 == "丘处机") | 筛选丘处机 |
💡 思路:
- 英雄和武器分别是图的两类顶点
- 英雄和武器之间的"使用关系"是边
- 用 triplets 视图拿到英雄名和武器名
- 按英雄聚合(reduceByKey)
- 筛选出丘处机
笔记版本:V1.0 对应教材:《Spark大数据技术与应用(第3版)》人民邮电出版社 对应项目:项目7 统计得分排名前10的网页——Spark GraphX图计算框架 最后更新:2026年8月