图计算 GraphX
概述
GraphX 是 Spark 上的图计算框架,把图表示为属性图(顶点+边携带属性),复用 Spark 的 RDD 基础设施做分布式图计算。本文讲清属性图模型、图算子、Pregel 迭代框架,以及 PageRank/TriangleCount 等内置算法与自定义算法。
一、属性图模型
1.1 定义
scala
// 图 = 顶点 RDD + 边 RDD
Graph[VD, ED](
vertices: RDD[(VertexId, VD)], // 顶点:id → 属性
edges: RDD[Edge[ED]] // 边:源、目标、属性
)| 部分 | 示例 |
|---|---|
| 顶点 | (1, "用户A"), (2, "用户B") |
| 边 | Edge(1, 2, "关注") |
| 属性 | 顶点/边上的任意数据 |
1.2 构建图
scala
val users: RDD[(Long, String)] = sc.parallelize(Seq(
(1L, "A"), (2L, "B"), (3L, "C")))
val edges: RDD[Edge[Int]] = sc.parallelize(Seq(
Edge(1L, 2L, 1), Edge(2L, 3L, 1)))
val graph = Graph(users, edges)1.3 图存储
| 布局 | 说明 |
|---|---|
| 顶点表 | 按 VertexId 分区 |
| 边表 | 按源/目标分区 |
| 路由表 | 边与顶点的对应关系 |
图计算通过多次迭代的消息传递收敛结果。
二、图算子
2.1 结构算子
| 算子 | 作用 |
|---|---|
| subgraph | 子图过滤(顶点+边条件) |
| reverse | 反转边方向 |
| mask | 求图交集 |
| groupEdges | 合并平行边 |
2.2 属性算子
| 算子 | 作用 |
|---|---|
| mapVertices | 变换顶点属性 |
| mapEdges | 变换边属性 |
| mapTriplets | 变换三元组(边+两端点) |
2.3 三元组
scala
graph.triplets.map(t =>
(t.srcAttr, t.dstAttr, t.attr) // 边 + 两端点属性
)三、Pregel API
3.1 思想
Pregel 是谷歌提出的顶点为中心迭代模型:每轮迭代,顶点接收邻居消息、更新自身状态、向邻居发消息,直到收敛。
迭代模型:
1. 初始消息发送给所有顶点
2. 每轮:顶点聚合收到的消息
3. 顶点更新状态(或保持不动)
4. 产生新消息发给邻居
5. 无消息产生 → 收敛结束3.2 API 签名
scala
graph.pregel(
initialMsg, // 初始消息
maxIterations, // 最大迭代数
activeDirection // 消息发送方向
)(
vprog, // 顶点程序:合并消息并更新顶点
sendMsg, // 发消息:根据边三元组决定消息
mergeMsg // 合并消息:多条消息合并
)3.3 示例:传播最小值
scala
val result = graph.pregel(
Long.MaxValue, 10, EdgeDirection.Out
)(
(id, attr, msg) => math.min(attr, msg),
triplet => {
if (triplet.srcAttr < triplet.dstAttr)
Iterator((triplet.dstId, triplet.srcAttr))
else Iterator.empty
},
(a, b) => math.min(a, b)
)四、内置算法
4.1 PageRank
衡量顶点重要性:迭代把自身权重分给邻居。
scala
val ranks = graph.pageRank(0.0001).vertices| 参数 | 说明 |
|---|---|
| tol | 收敛阈值 |
| resetProb | 随机跳转概率(默认 0.15) |
| 特点 | 说明 |
|---|---|
| 迭代式 | 每轮更新权重 |
| 收敛 | 权重变化小于阈值停止 |
| 适用 | 网页排序、影响力分析 |
4.2 ConnectedComponents 连通分量
scala
val cc = graph.connectedComponents().vertices| 说明 | 值 |
|---|---|
| 作用 | 找出图中相互可达的连通子图 |
| 原理 | 传播最小顶点 ID |
| 适用 | 社群发现、图划分 |
4.3 TriangleCount 三角计数
scala
val triangles = graph.triangleCount().vertices| 说明 | 值 |
|---|---|
| 作用 | 统计每个顶点参与的三角形数 |
| 适用 | 社群密度分析、社交网络 |
| 要求 | 图需三角排序(简化后) |
4.4 其他算法
| 算法 | 作用 |
|---|---|
| stronglyConnectedComponents | 强连通分量 |
| SVDPlusPlus | 协同过滤 |
| LabelPropagation | 标签传播(社群) |
| GraphOps | 度、入度、出度统计 |
五、自定义图算法
5.1 步骤
1. 用 graph.mapVertices/mapEdges 准备初始属性
2. 定义消息类型与发送逻辑
3. 用 pregel 迭代计算
4. 收敛后收集结果5.2 示例:无权图最短路径
scala
def shortestPath(graph: Graph[Long, Int], sourceId: Long) = {
graph.pregel(
Long.MaxValue, 100, EdgeDirection.Out
)(
(id, dist, newDist) => math.min(dist, newDist),
triplet => {
if (triplet.srcAttr + 1 < triplet.dstAttr)
Iterator((triplet.dstId, triplet.srcAttr + 1))
else Iterator.empty
},
(a, b) => math.min(a, b)
)
}| 设计要点 | 说明 |
|---|---|
| 状态语义 | 顶点属性即算法状态 |
| 消息最小化 | 只发变化的消息 |
| 收敛条件 | 无新消息或达迭代上限 |
六、性能与限制
6.1 性能要点
| 要点 | 说明 |
|---|---|
| 迭代代价 | 每轮 Shuffle,减少迭代次数 |
| 消息量 | 消息越少收敛越快 |
| 分区 | 按边源顶点分区,本地性 |
| 内存 | 图数据 + 消息内存开销 |
6.2 限制
| 限制 | 说明 |
|---|---|
| 不适合超大图迭代 | 每轮全量 Shuffle |
| 图灵完备性 | 强于 SQL 弱于专用图库 |
| 对比 Neo4j | 分布式批处理 vs OLTP 图查询 |
| 选型 | 场景 |
|---|---|
| GraphX | 批式图计算(PageRank 等) |
| Neo4j | 在线图查询 |
| GraphFrames | DataFrame 之上的图 API(更易用) |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 迭代太慢 | 消息量过大,优化消息逻辑或减少迭代 |
| 内存 OOM | 图数据过大,增加分区与内存 |
| 结果不收敛 | 检查消息发送条件与收敛阈值 |
| 顶点 ID 冲突 | 统一为 Long 类型 VertexId |
| 需要 SQL 风格操作 | 用 GraphFrames 或转 DataFrame |