Apache Spark — 慢作业的答案在执行计划和事件日志里
explain 是一张收据,显示 Spark 决定少读了什么
一句话总结
Spark 不会把你写的代码原样运行。优化器会修改计划:把条件下推给文件读取器、裁掉用不到的列、跳过不需要的分区目录。结果会以 PushedFilters、ReadSchema、PartitionFilters 的形式打印在 explain 的物理计划里,而在运行过程中由 AQE 再改一次的最终计划,要等到行动算子结束之后才能看到。
为什么要读计划
两个结果相同的查询,一个 3 秒,一个 3 分钟。代码看上去几乎一样。差别通常出在读了多少。比如把一个条件套上函数,结果读了整个 Parquet 文件;或者因为 select("*") 这一行,把不需要的二十列从磁盘上载了上来。
这样的差别,代码读得再仔细也看不出来。要看 Spark 把代码转换成了什么样的计划。那张收据就是 explain。这里有一种重要的态度:不要相信优化,要确认。优化器只在做得到的时候才去做,没做到的时候不会发出警告。
工作原理
根据 SQL EXPLAIN 参考,EXTENDED 会展示四种计划。从查询中提取出的解析后的逻辑计划、确定了名字和类型的分析后的逻辑计划、经过优化规则的优化后的逻辑计划,以及真正要运行的物理计划。PySpark 的 explain 用模式来选择这些。默认(simple)只有物理计划,extended 同时给出逻辑和物理计划,codegen 给出生成的代码,cost 在有统计信息时给出逻辑计划和统计信息,formatted 则把物理计划拆成概览和各节点的详情来展示。
比较优化后的逻辑计划和物理计划,就能看到优化器做了什么。大部分都集中在读取文件的那一行节点上。
+- FileScan parquet [order_id#0,qty#3,day#7]
PartitionFilters: [isnotnull(day#7), (day#7 = 2026-01-03)]
PushedFilters: [IsNotNull(qty), GreaterThanOrEqual(qty,5)]
ReadSchema: struct<order_id:string,qty:int>
条件下推(PushedFilters)。意思是 where(qty >= 5) 不只是作为 Filter 节点存在,还被交给了文件读取器。在 Parquet 配置中,spark.sql.parquet.filterPushdown 的默认值是 true。交下去的条件能跳过什么,取决于文件拥有的统计信息。利用统计信息一节列举了 Spark 直接从数据源读取的统计信息的例子:Parquet 元数据中的计数以及最小值、最大值。只看最小值、最大值,就知道无法满足条件的块不必打开。
列裁剪(ReadSchema)。只读取查询自始至终要用的列。Parquet 按列聚集存储,所以只读两列的话,其余列的字节不会从磁盘上载上来。CSV 是按行的,所以这个好处要小得多。把文件转换成 Parquet 一次,道理就在这里。
分区裁剪(PartitionFilters)。这是作用在值包含于目录名里的列(例如 day=2026-01-03/)上的条件。Parquet 的分区发现会从这样的路径中自动提取分区列并放入 schema。作用在这个列上的条件,在文件打开之前的目录列表阶段就处理完了,所以不相关的日期,成本为 0。
合并连续的过滤。即使像 where(a).where(b) 这样把条件拆开写多次,在优化后的逻辑计划中也会合并成一个 Filter。所以为了可读性而把条件拆开写,并不损害性能。
下推失效的瞬间
下推在条件是列本身与常量的比较时效果最好。读取器能听懂的,是“qty 大于等于 5”这样的简单形态。像 upper(status) = 'PAID' 这样把列套上函数,就不再是可以交给读取器的形态,该条件就会从 PushedFilters 里消失。结果相同,没有错误,也没有警告。只是读取量增加。本实验中要用计划确认的正是这个差别。
分区列也是如此。在日期分区上加 day = '2026-01-03',会进入 PartitionFilters,但如果先对日期字符串做加工再过滤,裁剪可能不会发生。“过滤条件要加在保持存储形态的列上”这个习惯就是由此而来。
AQE——运行过程中改变的计划
根据性能调优文档,自适应查询执行(AQE)是在运行过程中利用运行时统计信息重新优化计划的技术,从 Spark 3.2.0 起默认开启。这句话对读计划有重要的影响:在行动算子之前打印的 explain 不是最终计划。
在行动算子之前打印,最上面会看到 AdaptiveSparkPlan isFinalPlan=false。要等 shuffle 阶段真正结束、知道了大小之后,AQE 才会合并 shuffle 分区或更换连接方式。调用一次行动算子后,再看同一个 DataFrame 的计划,就会看到 isFinalPlan=true,以及 AQEShuffleRead 这样的节点。要确认实际运行的计划,就得看这个最终计划。在 Web UI 文档的 SQL 选项卡里,也能对每个查询通过 Details 展开这四种计划,并查看每个算子通过了多少行之类的指标。
在现场相遇的样子
“明明加了过滤,为什么全都读了?” 最常见的原因是:套了函数的条件、类型不同的比较(字符串列与数字常量)、以及直接读取 CSV。这三种情况,都会在计划的 PushedFilters 和 ReadSchema 那一行里暴露出来。
“select 是后面才做的,所以没关系”,大体上是对的。因为优化器只会留下自始至终要用的列。不过,如果中途做了 cache(),或者把整行传给 Python UDF,在那个位置需要的列就会增多。要始终用 ReadSchema 确认。
没有统计信息,计划就会摇摆。统计信息一节写道,统计信息缺失或不准确时,Spark 无法选出好的计划,并提示用 explain(mode="cost") 查看估算值,用 SQL UI 的 isRuntime=true 查看运行中的统计信息。
实际工作中真正重要的事
- 不要相信优化,用 explain 确认。没做到的时候不会发出警告。
- 读取节点的那一行是关键。先看 PushedFilters、ReadSchema、PartitionFilters 这三栏。
- 条件要加在保持存储形态的列上。套上函数,下推就会悄悄消失。
- 分区列上的条件会跳过目录。这是最便宜的裁剪,所以把经常用来过滤的列设为分区键。
- AQE 的最终计划要在行动算子之后看。isFinalPlan=false 的计划只是预告片。
下一项实验要做什么
把订单 CSV 转换成加上日期列的 Parquet,再打印 explain 的 extended 和 formatted 计划。确认条件和列选择在 PushedFilters 和 ReadSchema 里是怎么显示的,并比较同一个条件套上 upper() 之后下推消失的情况。在按日期拆分写出的 Parquet 上选出某一天,看条件显示在 PartitionFilters 里;把按客户的聚合作为行动算子运行之后,读取 AQE 最终计划和事件日志,数一数在 200 个 shuffle 分区中,实际读取了 shuffle 的任务有几个。最后,在优化后的逻辑计划里,确认连续的过滤被合并成一个,1 + 2 这样的常量表达式被预先算好。