ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

航班飞行网图分析实战:基于Spark GraphX的图计算应用

航班飞行网图分析实战:基于Spark GraphX的图计算应用 做航班数据分析这几年我越来越觉得把航班数据当普通关系表来算其实有点浪费。航线的本质就是一张巨大的图——机场是顶点航班是边旅客中转天然就是在图上做路径遍历。刚接触 Spark 的时候我也习惯性用 DataFrame 做 join 和 groupBy但遇到某个机场被哪些枢纽串联哪些区域因为一个枢纽断掉就彻底失联这类问题SQL 写起来绕且慢。后来我终于转到 GraphX 上系统做了一套航班飞行网图分析这应该算是我在 Spark 生态里回报率最高的一次技术投入。这篇文章就从项目实战的角度把整个航班飞行网图分析的过程拆开讲清楚包括为什么非要用图、怎么把原始航班数据构建成 GraphX 的顶点和边、PageRank 和连通分量这些算法在航班场景里怎么落地、以及跑大规模图任务时内存分区那些坑。内容面向已经会写 Spark SQL、但对 GraphX 比较陌生的读者也适合做数据分析和数据挖掘的朋友参考。1. 为什么航班网络分析最后落到了图计算上先说我之前用关系模型处理航班数据遇到的具体瓶颈。假设有一个航班表里面有起飞机场、落地机场、航班号、日期、机型、座位数这些字段。要回答哪些机场是整个网络的枢纽这个问题关系模型的思路是先统计每个机场的起降架次再计算吞吐量然后看排名前 20。这个其实并不难一条 groupBy 就能搞定。但要是问题变成从北京出发最多中转两次能到达哪些欧洲机场SQL 就得自连接两到三次数据量大时中间结果膨胀得非常快。更麻烦的是如果要求中转次数不固定比如任意多次中转找出所有能从北京到达的机场这就是一个递归遍历问题纯 SQL 表达起来相当痛苦性能和可读性都很难兼顾。图模型在这里就顺手得多。把机场抽象成顶点把两座机场之间的直飞航线抽象成有向边航班问题的语义几乎一比一映射到图论里。整个网络的拓扑结构就是一张有向图 G (V, E)V 是机场集合E 是航线集合。基于这张图上面那一堆业务问题就变成几个标准的图算法问题哪个机场最重要 —— PageRank 或中心性指标从某地出发能到哪里 —— 可达性问题也就是连通分量哪些机场构成紧密联系的小团体 —— 三角计数或社区发现哪个机场是全网连接的关键 —— 割点和桥接边航班的直接接驳关系 —— 一阶或二阶邻居查询GraphX 的价值不在于提供了多少花哨的算法而在于它把这些图算法做成可分布式、可横向扩展的算子。航班数据量级大的时候几十万顶点、几百万条边单机用 NetworkX 可能还没崩但一旦涉及复杂遍历和迭代内存就成了瓶颈。GraphX 底层有 RDD 的分区机制撑着数据可以分散在集群里计算每条边和顶点被打散到不同的分区迭代时尽量不 shuffle 掉所有数据这在生产环境里是核心优势。2. 从 CSV 到 Graph构建 verticesRDD 和 edgesRDD 的完整链路GraphX 里最核心的抽象就是 Graph它由两个 RDD 组成VertexRDD 和 EdgeRDD。VertexRDD 的格式是(VertexId, VD)VertexId 在图上所有顶点里必须唯一VD 是顶点属性EdgeRDD 的格式是(srcId, dstId, ED)srcId 是边的起点dstId 是终点ED 是边属性。所以整个项目的第一步就是把原始航班数据映射成这两种结构。这个过程看着简单里面细节不少。我用的是国内某数据平台导出的航班计划数据大致字段如下字段示例值含义flight_noCA1801航班号origin_airportPEK起飞机场三字码dest_airportSHA目的机场三字码dep_time08:00计划起飞时间arr_time10:15计划到达时间aircraft_typeA321机型seats186可用座位数frequency1111111一周内执飞1 表示执飞构建顶点时要注意顶点 ID 不能直接用字符串 PEKGraphX 的 VertexId 是 Long 类型。所以先要对机场三字码做编码。我当时用的是简单的字符串哈希加碰撞处理更稳妥的做法是维护一个机场字典表给每个三字码分配一个从 1 开始的递增 ID。这一步最好在 DataFrame 阶段就完成我直接开了个monotonically_increasing_id()来兜底。顶点属性我放的是一个自定义的 AirportInfo case class包含三字码、城市名、机场名。以前图构建的顶点属性存的是原始字符串后续做分析时每次都要反复解码后来发现直接在顶点属性里把常用信息都塞进去后续聚合展示就省了反复查表。边的构建相对更细。一条原始航班记录不一定只生成一条边。如果两台机场之间一天有多个航班我是在边上把航班列表聚合进边属性而不是每条航班记录生成一条边。这样图就不会因为航班量放大而爆炸也更符合航线网络分析而不是航班时刻表分析的粒度。边属性我定义成 RouteInfo里面存的就是航班号列表、总班次、机型集合。val flightsDF spark.read .option(header, true) .option(inferSchema, true) .csv(/data/flights/2024_full.csv) val airportDim flightsDF .select(col(origin_airport), col(origin_city), col(origin_name)) .distinct() .collect() .zipWithIndex .map { case (row, idx) val code row.getString(0) (idx.toLong, AirportInfo(code, row.getString(1), row.getString(2))) } val vertexRDD: RDD[(VertexId, AirportInfo)] spark.sparkContext.parallelize(airportDim.toSeq) val edgeRDD: RDD[Edge[RouteInfo]] flightsDF .groupBy(origin_airport, dest_airport) .agg( collect_list(flight_no).as(flights), count(flight_no).as(flight_count), collect_set(aircraft_type).as(aircraft_types) ) .rdd .map { row val srcCode row.getString(0) val dstCode row.getString(1) val srcId codeToVertexId(srcCode) val dstId codeToVertexId(dstCode) Edge(srcId, dstId, RouteInfo(row.getAs[Seq[String]](flights), row.getLong(flight_count), row.getAs[Seq[String]](aircraft_types))) } val graph Graph(vertexRDD, edgeRDD)这段代码里codeToVertexId需要先把三字码到 ID 的映射广播出去或者用 map 查字典。千万别说你直接在 map 里 collect 再查Driver 端一个大集合反复 broadcast 到各 executor任务多了会拖垮网络。我是把机场编码表做成 Broadcast 变量构建 edges 时在 executor 端直接查性能好很多。3. 跑通四个核心算法度数、PageRank、连通分量、三角计数实战图建好之后真正的分析才开始。我在这套航班项目里用到的核心算法主要是四个每个解决一类业务问题组合起来就能把整个航班网络的全貌拼出来。3.1 入度、出度和总度机场繁忙程度的第一层判断度数是图论里最基础的指标。在航班网络里出度表示从这个机场直飞能到达多少个不同的机场入度表示有多少个机场直飞能到达这里总度则是这个机场在整个航线网络中的直接连接广度。GraphX 里degrees、inDegrees、outDegrees都是预置算子直接调用就行。需要注意的是GraphX 的degrees默认把每条无向边算一次对于有向图来说你算总度时要把入度和出度分别算清楚再合并否则航线方向的信息就丢了。val inDeg graph.inDegrees val outDeg graph.outDegrees val degreeStats inDeg.fullOuterJoin(outDeg) .map { case (vid, (inOpt, outOpt)) val in inOpt.getOrElse(0) val out outOpt.getOrElse(0) (vid, in, out, in out) } .sortBy(_._4, ascending false)在实际数据上跑出来的结果很有说服力。全网连接边数最多的那几个机场基本就是行业里默认的三大门户机场和几个区域性枢纽。但有一个点很容易被忽略出度和入度在网络里大部分时候是不相等的因为有些机场只有进港航班没有出港航班比如某些以旅游目的地为主、返程航班被识别成另一条航线的情况这个不对称在后续做可达性分析时很关键。3.2 PageRank从谁连接多到谁连接得重要度数的最大问题是它把所有连接一视同仁。一个机场连接了 10 个小机场跟一个机场连接了 10 个大枢纽度数一样但后者的实际网络地位高得多。PageRank 恰好能解决这个问题一个节点的重要性取决于谁在指向它以及指向它的节点自身重不重要。GraphX 的 PageRank 实现是迭代式的参数tol控制收敛阈值。我实测下来tol设 0.01 和 0.001 对最终排名的影响很微小但迭代次数差了不少。航班场景里排名靠前的机场很稳定所以不需要为了精确度损失太多时间。val ranks graph.pageRank(0.0001).vertices val rankedAirports ranks.join(vertexRDD) .map { case (_, (rank, info)) (info.airportCode, rank) } .sortBy(_._2, ascending false)有个容易踩的细节PageRank 默认模型里存在随机跳转因子0.85这个参数在航班网络里的含义可以理解为乘客偶尔坐一次非枢纽连接航线的概率。保持默认即可——真去调它对排名结果的影响通常都小得可怜。从结果看PageRank 排名和按吞吐量的官方排名不完全一致这一点很有意思。有几个吞吐量很大的机场 PageRank 反而排不到前三原因是它们的流量高度集中在少数几条高频航线上连接到其他机场的种类不够多而一些区域枢纽因为连接了大量支线机场PageRank 反而上去了。这其实反映了两种不同类型的枢纽一种是深度型枢纽靠干线高频一种是广度型枢纽靠支线覆盖。这两类机场在整个网络里承担的角色不一样后续做航线布局和运力分配时需要区别对待。3.3 连通分量回答哪些机场和外界是断开的连通分量在航班网络里的语义很有意思。强连通分量表示两两之间都能通过一系列航线互相到达弱连通分量则忽略方向只关心是否有路径连在一起。GraphX 提供的connectedComponents算法返回的是弱连通分量实现原理是基于 Pregel 迭代传播最小顶点 ID直到每个顶点的分量标签不再变化。代码一行但理解它的输出很关键。每个顶点会得到一个分量 ID同一个分量里的顶点代表从拓扑上看它们属于同一个可互达的片区。我在真实数据上的应用是先算弱连通分量把机场按分量分组然后找出那些只包含一两个机场的小分量。这些小分量往往是数据质量问题——有些机场三字码在新版航班表里已被废弃或者某些包机航线只在特定季节运行平时根本没航班。把这些孤立分量筛出来反过来帮我们治理了源数据的脏数据。val cc graph.connectedComponents().vertices val ccStats cc.join(vertexRDD) .map { case (_, (componentId, info)) (componentId, info.airportCode) } .groupByKey() .map { case (componentId, airports) (componentId, airports.toList, airports.size) } .sortBy(_._3, ascending false)强连通分量更有业务价值。用stronglyConnectedComponents时要注意它对迭代次数的参数numIter极其敏感设小了结果不对设大了非常慢。我是在一个 500 个机场、3000 条边的图上跑的设 50 次迭代大约跑了十几秒还能接受。结果里最大的那个强连通分量基本就是全国或者说全网航线的核心骨架分量外的机场和主网的联系往往依赖某几条特定航线一旦这些航线断了它们就被隔离了。3.4 三角计数发现区域小团体和航线冗余度三角形在网络里的含义是三个节点两两相连。在航班网络里A-B-C 三条航线如果都存在说明这三个城市之间存在比较密集的往来旅客往返可以有更灵活的组合方式。三角计数高的机场往往属于某个内部联系紧密的区域集群。GraphX 的triangleCount要求边是无向的或者至少传进去的图要满足一定条件。它要求数据按srcId dstId的形式组织否则结果会出错。我第一次跑的时候忽略了这一点直接拿有向图去跑跑出来的三角形数量明显不对后来查文档才发现有这个前置要求。val undirectedGraph graph .subgraph(epred edge edge.srcId edge.dstId) val triCounts undirectedGraph.triangleCount().vertices把三角计数结果和 PageRank 排名放一起分析能发现一类有趣的现象有些机场 PageRank 不高三角计数却很大说明它们在某个区域内和周边的连接非常充分是整个区域网络的组织者相反有些大枢纽三角计数反而不高因为它们大量连接的是点对点的远距离航线而不是区域内部的短途密集网。这两种角色的机场放在航线网络优化里的策略是完全不同的。4. 内存、分区和序列化GraphX 性能调优的关键GraphX 在 Spark 生态里从来不是性能最好的图计算引擎——GraphFrames、甚至专门的图数据库在某些场景下都可能超过它。但它胜在能直接嵌入现有的 Spark 数仓链路不用额外引入一套存储和计算引擎。在项目里跑了大几百万条边以后我总结出几个真正影响性能的因素按影响从大到小排。4.1 顶点 ID 的连续性决定了聚合效率GraphX 底层的 VertexRDD 使用哈希和索引来加速顶点查找顶点 ID 映射越紧凑、范围越小内部索引的存储和查找效率越高。在构建顶点时别用字符串哈希的 Long 值把 ID 空间撑得很大尽量用 1 到 N 的连续数值。我一开始偷懒直接对三字码做 MD5 再去截取结果顶点 ID 跨度巨大使用outerJoinVertices时 Shuffle 的数据量比连续 ID 大了将近一倍。后来改成按机场字典递增编号性能立刻好转。4.2 边分区策略如何减少迭代中的跨节点通信GraphX 默认的边分区策略是RandomVertexCut它会尽量把同一条边的两个端点相关的信息放到同一个分区里。对于 PageRank 这种迭代算法每轮迭代都需要在邻居间传递消息消息传递的通信开销基本由边分区方式决定。如果边的分区剪得不好大量消息需要跨 executor 传输网络带宽直接变成瓶颈。GraphX 里可以通过partitionBy重设分区策略val partitionedGraph graph .partitionBy(PartitionStrategy.EdgePartition2D, numPartitions 200)EdgePartition2D在大多数场景是比默认策略更稳的选择它在二维空间上把边切成网格让顶点在多个分区里的副本数尽量均衡。分区数怎么选我个人的经验公式是设为 executor 总数乘以每个 executor 的核心数再乘以 2 到 3。实际调整时观察 Spark UI 中每个 stage 的时间如果某个 stage 有严重的数据倾斜某个 task 运行时间比中位数长很多优先考虑增大分区数而不是改分区策略。4.3 迭代算法里的缓存和持久化级别PageRank、连通分量这类迭代算法有个共性每次迭代都要重复读取图的拓扑结构。如果你连续跑多个算法比如先 PageRank 再连通分量中间只要action被触发一次整个图就会重新计算。我习惯在第一个图操作之后立刻persist存储级别选MEMORY_AND_DISK。val cachedGraph graph .partitionBy(PartitionStrategy.EdgePartition2D, 200) .persist(StorageLevel.MEMORY_AND_DISK)persist之后记得在你的 Spark 应用结束时unpersist否则 GraphX 的缓存会一直占用 executor 内存影响后续其他作业的稳定性。我在开发时曾因为一个persist忘记释放导致集群上后续几个任务频繁 OOM排查半天才找到原因。4.4 顶点属性别塞太多复杂对象前面提到我在顶点属性里放了AirportInfo这是有代价的。每次迭代不管用得着用不着这些属性数据都会在序列化和反序列化过程中经过网络和磁盘。如果属性是一个嵌套很深的 case class序列化开销直线上升。我的建议是顶点属性只放这一步分析必需的字段比如在做 PageRank 分析时顶点属性只需要机场编码和名称其他城市、机型那些信息可以先单独建 DataFrame最后分析完再 join 回去。这个优化在数据量大的时候效果非常明显。5. 综合结果解读从图计算输出到业务决策算法跑完只是第一步把图计算结果翻译成业务语言才是项目的最终目的。这里就说几个我在实际项目里做过、并且确实被业务采纳的分析视角。5.1 枢纽机场的层次识别把度数、PageRank、连通度放一起看单个指标容易以偏概全。我把每个机场的度数、PageRank、所在连通分量大小、三角计数拼接成一张宽表然后用聚类的方式把机场分成几类。分出来的结果大致是这几类全球/全国级枢纽PageRank 和度数双高连接范围覆盖全国甚至洲际连通分量也是最大那个核心分量中的一员。区域枢纽PageRank 中等偏上度数较高三角计数很高是区域内小机场连接外部的主要跳板。深度干线节点PageRank 很高但度数不高典型就是那些靠几条航线高频撑起吞吐量的机场。支线末梢度数低PageRank 低三角计数接近零在整个网络里基本处于从属位置。这个分类后来直接用于航线补贴政策的制定对区域枢纽类机场重点加密它到全国枢纽的航线对支线末梢则优先补贴到临近区域枢纽的航线而不是盲目开通远程航线。5.2 网络脆弱性分析去掉某些机场整个网络会变成什么样连通分量的另一个高阶玩法是删点测试。我遍历 PageRank 排名前 20 的枢纽每次删掉一个顶点重新计算全图的最大连通分量大小看减少了多少。某两个枢纽被删除后最大连通分量急剧缩小说明它们是连接南北或者连接东西的关键桥接点。这类机场一旦出现长时间停运对整个网络的打击是结构性的不是简单削减运力能缓解的。这种脆弱性分析还牵引出一个现实问题当某个大枢纽因天气、流控等特殊原因大面积取消航班时通过图网络快速找到替代中转点让旅客在最少的额外中转次数内到达目的地这个需求后来发展成了一个 ODSOrigin-Destination-Segment替代路径推荐功能底层用的就是 GraphX 的 BFS 变体和最短路径思路。5.3 航线新增的模拟打分加一条边图和之前有什么变化最后分享一个比较有意思的玩法。我在已有图的基础上模拟添加一条候选新航线也就是加一条边然后重新算受影响顶点的 PageRank 和全图三角计数通过前后对比来判断这条新航线对整个网络的提升价值。这个做法本质上不是严谨的图论实验因为加一条边对局部顶点的影响远大于全局但作为快速业务判断工具非常有价值。实际操作时做一个待选航线列表生成新图批量计算指标变化量筛选出那些对目标区域连接度提升最大的航线候选。整个过程不需要复杂的因果推断模型图结构的变化本身就是一种信号。6. 踩过的坑GraphX 项目里最致命的几个细节最后把我在项目过程中遇到的、值得单独拎出来提醒的细节和坑集中总结一下。6.1 子图操作后顶点和边会不一致GraphX 的subgraph可以分别用vpred和epred过滤顶点和边。问题在于如果你只过滤了边那么某些顶点可能变成孤立点反过来如果你只过滤了顶点那么被删掉顶点的边依然存在于 EdgeRDD 里造成悬挂边。我在做脆弱性分析时第一次犯了这个错删掉一个枢纽机场的顶点但没过滤与之相连的边结果连通分量计算直接把那个被删掉的枢纽又算回去了前几名的结果显示几乎没有变化我差点下了网络很稳健的错误结论。正确做法是删点时同时用vpred和epred把相关边也过滤掉def removeAirport(g: Graph[AirportInfo, RouteInfo], targetId: VertexId) { g.subgraph( vpred (vid, _) vid ! targetId, epred e e.srcId ! targetId e.dstId ! targetId ) }6.2 字符串三字码在 ID 里的隐藏坑三字码有大小写区别有的源系统导出全大写有的混着大小写。如果构建顶点和边时用的三字码大小写不一致同一个机场会被拆成两个顶点。别笑这是真实发生过的。我当时在验证时发现全网机场数量竟然比官方数量多了两百多个最后定位到问题就出在一个上游表的城市名是区分大小写存的。所以构建顶点前务必统一格式.upper()一下再处理能省掉后面大量返工。6.3connectedComponents在大图上要小心 Driver 端 OOM代码connectedComponents().vertices执行完如果直接.collect()的话所有顶点的分量结果会全部回收到 Driver 端顶点数量如果是几百万这本身不是问题但如果顶点属性里塞了复杂的对象比如长字符串列表collect 出来的数据量会非常可观Driver 内存直接被撑爆。建议在 collect 之前先尽快 join 上所需的最小字段select 出需要展示或分析的列再回收数据。另外能用saveAsTextFile落盘的就不要全程走 Driver 内存。6.4 图的度、PageRank 之类的操作结果需要 rejoin 才能看到属性GraphX 返回的vertices是RDD[(VertexId, Double)]这种形态只有顶点 ID 和分数没有你想要的机场名和城市信息。初学者容易在拿到结果后一头雾水。记得在展示前把原始顶点属性 join 回来。这个 join 用 RDD 的 join 就行但要确保两边的 key 类型一致都是 Long不然又会引入意外错误。val resultWithName ranks.join(vertexRDD) .map { case (vid, (rank, info)) (info.airportCode, rank) }6.5 避免把图对象重复序列化传给 UDF有一种错误做法是在 GraphX 计算完成后把图对象作为广播变量传出去然后在 DataFrame 的 UDF 里反复查图索引。图对象如果不大倒还好说一旦图比较大广播到 executor 的副本本身就是巨大的内存负担而且 UDF 里对图做点查询可能要遍历边集合性能极差。正确做法是先在图计算阶段把需要的结果集抽出来转成普通的 Map 或 DataFrame再用于后续逻辑。7. 实战总结飞行网图分析项目带来的三个思考这趟实战做完我对 GraphX 的看法比刚开始时要务实不少。第一GraphX 不是说有了它就不需要 Spark SQL。事实上我的项目流程里大部分数据清洗和聚合都是用 DataFrame 完成的图计算只是处理关系结构那一层特定逻換。最舒服的组合是DataFrame 负责 ETLGraphX 负责网络分析最后把结果再回写成 DataFrame 供报表和下游使用。第二算法选型要从业务问题出发不是图算法听起来高级就用。比如找出全网最重要的五个机场PageRank 最合适找到从 X 出发能到达的所有机场连通分量最直接两个城市之间最快中转方案就得走最短路径或者 BFS。每个算法有它的适用边界别套模板。第三图分析做出来的东西要真正被业务用起来得在结果可解释性上下功夫。我在输出机场分类时给每个类别都配了一个业务侧的说法——什么叫区域枢纽、什么叫支线末梢管理层能看懂。如果只给一张 PageRank 分数列表再专业也推不下去。最后想说航班飞行网图这个题目其实很适合当作 GraphX 的入门实战项目。数据量适中、图结构语义清晰、算法输出容易验证。你不需要一个几十亿边的大图才能体会 GraphX 的价值几千条边、几百个顶点的图已经能带你完整走一遍从建图到算法调优再到结果解读的全过程。走完这一遍后面再遇到社交网络分析、供应链路径优化、资金流转图这类类似场景你基本就有一个现成的思路框架可以迁移过去了。
返回列表