5.2 核心API与高级操作:转换与消息的搬运账单


5.2 核心API与高级操作:转换与消息的搬运账单

本节摘要:GraphX 的算子分三类——属性转换(mapVertices/mapEdges,零搬运)、结构变换(subgraph/reverse,分区内过滤)、消息传递(aggregateMessages,有搬运)。本节逐类巡检它们的执行形态,并用一个协作网络案例把三类串起来。

上一节把图拆成了三张表,本节开始在三张表上"做功"。判别任何 GraphX 算子的第一步,是问一句:这个操作要不要让分属不同分区的数据碰面?答案直接决定 Stage 里有没有 Shuffle 边界。

第一类:属性转换,分区内的流水线

// 把每个顶点的人名换成大写:每个分区独立完成,纯窄依赖 val g2 = graph.mapVertices((_, name) => name.toUpperCase) // 边属性换算协作强度:同样不跨分区 val g3 = g2.mapEdges(e => e.attr * 2) // mapTriplets 能同时看见两端属性,但依旧零搬运 // 因为三元组对齐用的就是本分区内的边表与点副本 val g4 = g3.mapTriplets(t => if (t.srcAttr.length > 2) t.attr + 1 else t.attr)

这类算子引擎视角最便宜:每个 Executor 线程扫自己那份分区,输出新的 RDD 血统,不产生任何网络流量。血统照旧惰性,只有行动算子触发时才真正执行。

第二类:结构变换,过滤为主

// 只留协作超过 5 次的边:边表过滤,分区独立 val strong = graph.subgraph(epred = e => e.attr > 5) // 只留"还在职"的点和相关边:点边谓词各管一摊 val active = graph.subgraph( vpred = (id, name) => name != "钱工") strong.numEdges // 4 条边剩 3 条:12、7、9 保留,3 被滤掉 active.numEdges // 钱工的点被删,3到4 的边随之消失,剩 3 条

subgraph 值得注意的细节:点谓词过滤后,悬空的边会被自动清理——边依赖点,物理上边表里那条边的目标点属性副本被剔除,结构上等价于删边。reverse 则只交换每条边的方向描述,同样是分区内的轻操作。

⚠️ 常见坑:subgraph 之后顶点副本的冗余不会立刻收缩,路由表也还是旧的。引擎在下次需要对齐时才用新图重新构建。若过滤后图要长期迭代,先做一次带缓存的重建(Graph(graph.vertices.filter(...), ...))再迭代更干净。

第三类:消息传递,图计算的心脏

aggregateMessages 是 GraphX 唯一的原生消息原语,一切内置算法都建筑在它之上。

import org.apache.spark.graphx.util.GraphGenerators // 场景:给每个人统计"直接协作过的总次数" // sendMsg:沿边发消息,三元组里能看见两端属性 // mergeMsg:同一顶点收到多条消息时的合并函数 val coop: VertexRDD[Int] = graph.aggregateMessages( ctx => { ctx.sendToSrc(ctx.edge.attr) // 给源点发:这次协作计 attr 次 ctx.sendToDst(ctx.edge.attr) // 给目标点也发一份 }, (a, b) => a + b // 收到多条就累加 ) coop.collect().foreach(println) // (4,9) 钱工只有一条边 // (1,15) 李工:12 + 3 // (2,19) 王工:12 + 7 // (3,19) 赵工:7 + 3 + 9

引擎视角的执行过程分三段。第一段 sendMsg 在边分区内完成,产出 (目标顶点id, 消息) 形式的中间 RDD——这一步零搬运,因为边和源点副本同分区。第二段是隐式的 Shuffle:消息按目标顶点的分区哈希重新分布,这是一次真实的网络搬运,量级等于"发出的消息总条数 × 消息大小"。第三段 mergeMsg 在目标分区内把同一顶点的消息折叠。三段正好对应窄依赖 → Shuffle 边界 → 窄依赖,两个 Stage。

把消息量控制住,是图算法调优的第一要务。ctx.sendToSrc 与 sendToDst 按需选用,只给一端发消息,搬运直接减半。

joinVertices:把结果贴回图

aggregateMessages 的产物是裸的 VertexRDD,通常要贴回图进入下一轮。

// 把协作度贴回顶点属性,类型变了就先换元组 val enriched = graph.outerJoinVertices(coop) { (_, name, deg) => (name, deg.getOrElse(0)) } enriched.vertices.collect().foreach(println) // (2,(王工,19)) (4,(钱工,9)) ...

outerJoinVertices 底层借助路由表做定向对齐:先广播小表或按位图定位,避免全局 Shuffle。左外语义保证没收到消息的点(比如孤点)不被丢掉,getOrElse 兜底。

巡检案例:一条时间线看三类算子的 Stage 形态

背景:运维要给协作网络做体检——过滤弱连接,给剩下的人标注协作强度。操作分三步:subgraph 过滤、aggregateMessages 统计、outerJoinVertices 回贴。结果在 Spark UI 上:subgraph 与后面的统计被 Catalyst 式地合并进同一个 Stage 的流水线,只有 aggregateMessages 的消息重分区形成唯一的 Shuffle 边界,整个作业两个 Stage。解读:三类算子混用时,决定 Stage 数量的只有消息传递类操作。变式:若把统计结果 collect 回 Driver 再广播贴回,小图可行,大图会撞 Driver 内存上限——outerJoinVertices 的分布式对齐正是为后者准备的。

本节要点回顾

  • 三类算子三档成本:属性转换零搬运、结构变换分区过滤、消息传递付一次 Shuffle
  • aggregateMessages 三段式:区内发消息、按目标重分区、区内合并,天然两个 Stage
  • 消息量是调优旋钮:只给一端发、过滤后再发、合并函数尽早折叠
  • outerJoinVertices 走路由表:定向对齐替代全局 Shuffle,左外语义保留孤立顶点
  • subgraph 的清理是逻辑等价:物理副本冗余要到下次对齐才收缩

下一节把这套心脏装进三个内置算法,逐轮拆它们的迭代账单。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U