做航班数据分析这几年,我越来越觉得把航班数据当普通关系表来算,其实有点浪费。航线的本质就是一张巨大的图——机场是顶点,航班是边,旅客中转天然就是在图上做路径遍历。刚接触 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_no | CA1801 | 航班号 |
| origin_airport | PEK | 起飞机场三字码 |
| dest_airport | SHA | 目的机场三字码 |
| dep_time | 08:00 | 计划起飞时间 |
| arr_time | 10:15 | 计划到达时间 |
| aircraft_type | A321 | 机型 |
| seats | 186 | 可用座位数 |
| frequency | 1111111 | 一周内执飞,1 表示执飞 |
构建顶点时要注意,顶点 ID 不能直接用字符串 "PEK",GraphX 的 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 的枢纽,每次删掉一个顶点,重新计算全图的最大连通分量大小,看减少了多少。某两个枢纽被删除后,最大连通分量急剧缩小,说明它们是连接南北或者连接东西的关键桥接点。这类机场一旦出现长时间停运,对整个网络的打击是结构性的,不是简单削减运力能缓解的。
这种脆弱性分析还牵引出一个现实问题:当某个大枢纽因天气、流控等特殊原因大面积取消航班时,通过图网络快速找到替代中转点,让旅客在最少的额外中转次数内到达目的地,这个需求后来发展成了一个 ODS(Origin-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 负责 ETL,GraphX 负责网络分析,最后把结果再回写成 DataFrame 供报表和下游使用。
第二,算法选型要从业务问题出发,不是图算法听起来高级就用。比如"找出全网最重要的五个机场",PageRank 最合适;"找到从 X 出发能到达的所有机场",连通分量最直接;"两个城市之间最快中转方案",就得走最短路径或者 BFS。每个算法有它的适用边界,别套模板。
第三,图分析做出来的东西,要真正被业务用起来,得在结果可解释性上下功夫。我在输出机场分类时,给每个类别都配了一个业务侧的说法——什么叫"区域枢纽"、什么叫"支线末梢",管理层能看懂。如果只给一张 PageRank 分数列表,再专业也推不下去。
最后想说,航班飞行网图这个题目其实很适合当作 GraphX 的入门实战项目。数据量适中、图结构语义清晰、算法输出容易验证。你不需要一个几十亿边的大图才能体会 GraphX 的价值,几千条边、几百个顶点的图已经能带你完整走一遍从建图到算法调优再到结果解读的全过程。走完这一遍,后面再遇到社交网络分析、供应链路径优化、资金流转图这类类似场景,你基本就有一个现成的思路框架可以迁移过去了。