Apache Spark — 慢作业的答案在执行计划和事件日志里
用 explain 确认少读了多少 Parquet 数据
目标
把订单转换成 Parquet,通过执行计划来读懂:同一个问题,随着条件的写法不同,能少读多少文件(下推、列裁剪、分区裁剪)。并通过事件日志确认,自适应执行(AQE)运行之后把计划改成了什么样。
为什么重要
Spark 慢的时候,最先该看的不是代码,而是计划。结果相同的两行代码,一个只从文件里读取需要的列和行,另一个全部读完再丢掉。这个差别不会出现在结果里,只会出现在计划里。
像 Parquet 这样的列式格式可以做两件事。只读取需要的列(列裁剪,计划里的 ReadSchema),并利用行组的最小值、最大值统计跳过不符合条件的块(条件下推,计划里的 PushedFilters)。再加上按目录拆分的分区,用条件选出来,就可以整个目录地跳过(PartitionFilters)。这三者都只有在优化器能理解条件的时候才会发生。列一旦被套上函数,下推就会消失。
AQE 会在 shuffle 结束之后,根据实际大小修改剩余的计划。所以运行之前的计划(isFinalPlan=false)和运行之后的计划可能不同,真正运行了什么,要在运行之后的计划或事件日志里看。
步骤
- 用 /root/spk/plan/convert.py(应用
spk-plan-convert)按 schema 读取订单,加上order_date(日期)列,以 Parquet 写入 /root/spk/plan/orders_pq。 - 在辅助文件 /root/spk/plan/planhelp.py 中放入
explain、field、items函数,并用 /root/spk/plan/explain.py(应用spk-plan-explain),把选出退款订单的order_id、qty的查询的extended、formatted计划分别保存到 /root/spk/plan/out/extended.txt 和 /root/spk/plan/out/formatted.txt,并对该查询执行collect()。 - 用 /root/spk/plan/pushdown.py(应用
spk-plan-push),把status == 'refunded'且qty >= 5的订单的order_id、qty、channel以 CSV 写入 /root/spk/plan/out/pushdown,并把该计划的PushedFilters条目列表和ReadSchema写入 /root/spk/plan/out/pushdown.json。 - 用 /root/spk/plan/nopush.py(应用
spk-plan-nopush),把同一个问题改成upper(status) == 'REFUNDED',写入 /root/spk/plan/out/nopush,并把该计划的PushedFilters列表写入 /root/spk/plan/out/nopush.json。 - 用 /root/spk/plan/partition.py(应用
spk-plan-part),按order_date拆分写入 /root/spk/plan/orders_by_day,选出 2026-02-14 这一天,把数出的行数和计划的PartitionFilters写入 /root/spk/plan/out/partition.json。 - 用 /root/spk/plan/aqe.py(应用
spk-plan-aqe),对按客户统计的订单数执行collect()之后,把计划保存到 /root/spk/plan/out/aqe_final.txt,并从该应用的日志里数出读取了 shuffle 的任务数,以{"shuffle_partitions": 200, "reduce_tasks": 정수}(占位符为整数)写入 /root/spk/plan/out/aqe.json。 - 用 /root/spk/plan/optimize.py(应用
spk-plan-opt),把串联qty > 1和qty > 3两个where、并把1 + 2选作three的查询的extended计划保存到 /root/spk/plan/out/optimized.txt,并执行collect()。 - 在 /root/spk/plan/report.md 中以
## 밀어내기(韩文,意为“下推”)、## 가지치기(韩文,意为“裁剪”)、## AQE三节写成报告。第二节放入第 5 步的行数,第三节放入第 6 步的任务数,都用数字写出。
参考
- 脚本放在
/root/spk/plan,并在那里执行spark-submit。同一文件夹的planhelp.py用from planhelp import explain来使用。 explain的模式:simple(只有物理计划)、extended(解析、分析、优化后的逻辑计划 + 物理计划)、codegen、cost、formatted(物理计划树 + 各算子详情)。各算子的PushedFilters、ReadSchema、PartitionFilters在 formatted 中每个占一行。- 在事件日志里,读取了 shuffle 的任务,是
SparkListenerTaskEnd的Task Metrics→Shuffle Read Metrics→Total Records Read大于 0 的任务。 - 常见错误:在 CSV 里找下推(CSV 必须把文件读完);在
collect()之前打印计划,却当作 AQE 的最终计划;因为有PushedFilters就断定少读了文件(行组统计必须能把条件分开,才会跳过)。 - 官方文档:Performance Tuning · EXPLAIN · Parquet Files · Adaptive Query Execution
把 CSV 转换成 Parquet
以应用名 spk-plan-convert 创建 /root/spk/plan/convert.py,按 schema 读取 /data/shop/orders.csv,加上 order_date = to_date(order_ts) 列,以 Parquet 写入 /root/spk/plan/orders_pq(不做拆分)。
列式文件按列分别存储,并且每个行组都有最小值、最大值统计。这两点是后面所有步骤中“少读”的原料。评分器会检查行数、列列表,以及文件是否真的是 Parquet(前四个字节是 PAR1)。
explain 的模式——从逻辑计划到物理计划
在辅助文件 /root/spk/plan/planhelp.py 中放入:把 explain 输出作为字符串取回的 explain(df, mode),从 formatted 计划中取出 이름: 값(占位符依次为名称、值)这一行的值的 field(plan, name),以及把 [A(x,1), B(y)] 拆成条目列表的 items(text);并以应用名 spk-plan-explain 创建 /root/spk/plan/explain.py,对 orders_pq 中选出 status == 'refunded' 的行的 order_id、qty 的查询,把 extended、formatted 计划分别保存到 /root/spk/plan/out/extended.txt 和 /root/spk/plan/out/formatted.txt,然后对该查询执行 collect()。
extended 会依次打印四节——解析后的逻辑计划、分析后的逻辑计划、优化后的逻辑计划、物理计划。在分析阶段,列名会被绑定到实际的列(#编号),在优化阶段,条件会下沉到靠近扫描的位置。formatted 把那个物理计划按算子逐个展开。评分器会检查 formatted 文件是否与实际运行的计划(事件日志)相同。
条件下推与列裁剪
以应用名 spk-plan-push 创建 /root/spk/plan/pushdown.py,把 orders_pq 中 status == 'refunded' 且 qty >= 5 的行的 order_id、qty、channel 以带表头行的 CSV 写入 /root/spk/plan/out/pushdown,并把该查询 formatted 计划中 PushedFilters 的条目列表和 ReadSchema 的值,以 {"pushed_filters": [문자열…], "read_schema": 문자열}(占位符均为字符串)写入 /root/spk/plan/out/pushdown.json。
PushedFilters 是交给 Parquet 读取的条件,ReadSchema 是实际读取的列。看看没有选出的 customer_id、product_id 是如何从 ReadSchema 里消失的,而只用在条件里的 status 又是如何留下来的。评分器会检查你写的列表是否与实际运行的计划里的逐字相同。
套上函数,下推就会消失
以应用名 spk-plan-nopush 创建 /root/spk/plan/nopush.py,把与第 3 步相同的问题改成 F.upper("status") == "REFUNDED" 和 qty >= 5,以带表头行的 CSV 写入 /root/spk/plan/out/nopush,并把该计划的 PushedFilters 条目列表,以 {"pushed_filters": [문자열…]}(占位符为字符串)写入 /root/spk/plan/out/nopush.json。
结果的行和第 3 步完全一样。不同的是计划。Parquet 不知道 upper(status) 是什么,所以这个条件不会交给读取,而是在文件全部读完之后由 Spark 来过滤。在实际工作中,把日期列套上 date_format 再比较,经常会掉进这个陷阱。
分区裁剪——整个目录跳过
以应用名 spk-plan-part 创建 /root/spk/plan/partition.py,把 orders_pq 用 partitionBy("order_date") 写入 /root/spk/plan/orders_by_day,再读取并数出 order_date == '2026-02-14' 的行,把这个值和该查询计划的 PartitionFilters 值,以 {"rows": 정수, "partition_filters": 문자열}(占位符依次为整数、字符串)写入 /root/spk/plan/out/partition.json。
被拆分的列不在文件里面,而是在目录名(order_date=2026-02-14)里。作用在这个列上的条件用来选目录,没被选中的日期的文件连打开都不会打开。评分器还会检查事件日志里“读取的分区数”指标是不是 1。
AQE 运行之后修改的计划
以应用名 spk-plan-aqe 创建 /root/spk/plan/aqe.py,对 orders_pq 中按客户统计的订单数(groupBy("customer_id").count())执行 collect() 之后,把 explain(mode="simple") 的输出保存到 /root/spk/plan/out/aqe_final.txt。然后在该应用的事件日志里数出读取了 shuffle 的任务数,以 {"shuffle_partitions": 200, "reduce_tasks": 정수}(占位符为整数)写入 /root/spk/plan/out/aqe.json。
spark.sql.shuffle.partitions 是 200,但这份数据的 shuffle 只有几 MB,所以 AQE 会合并分区。运行之后的计划里会看到 isFinalPlan=true 和合并后的 shuffle 读取(AQEShuffleRead)。读取了 shuffle 的任务数,可以在日志的 SparkListenerTaskEnd 里数 Total Records Read 大于 0 的任务得到。
优化器合并和折叠的东西
以应用名 spk-plan-opt 创建 /root/spk/plan/optimize.py,把 orders_pq.where(qty > 1).where(qty > 3).select("order_id", (lit(1) + lit(2)).alias("three")) 的 extended 计划保存到 /root/spk/plan/out/optimized.txt,并对该查询执行 collect()。
分析后的逻辑计划里有两个 Filter,(1 + 2) 也原样在那里。在优化后的逻辑计划里,两个条件合并成一个 Filter(过滤合并),1 + 2 折叠成 3(常量折叠)。这一步是要看到:你写的代码的样子,和实际运行的样子是不同的。
用数字写出少读了多少
在 /root/spk/plan/report.md 中写出 ## 밀어내기(韩文,意为“下推”)、## 가지치기(韩文,意为“裁剪”)、## AQE 三节。第二节用数字放入第 5 步的行数,第三节用数字放入第 6 步的 shuffle 读取任务数。
第一节写第 3、4 步里 PushedFilters 有什么变化,第二节写通过分区裁剪跳过了什么,第三节写从 200 减少到了几个。