TT Lab
开始
学习 学习路径 课程

Apache Spark — 慢作业的答案在执行计划和事件日志里

用 explain 确认少读了多少 Parquet 数据

在 TT Lab 中继续学习

目标

把订单转换成 Parquet,通过执行计划来读懂:同一个问题,随着条件的写法不同,能少读多少文件(下推、列裁剪、分区裁剪)。并通过事件日志确认,自适应执行(AQE)运行之后把计划改成了什么样。

为什么重要

Spark 慢的时候,最先该看的不是代码,而是计划。结果相同的两行代码,一个只从文件里读取需要的列和行,另一个全部读完再丢掉。这个差别不会出现在结果里,只会出现在计划里。 像 Parquet 这样的列式格式可以做两件事。只读取需要的列(列裁剪,计划里的 ReadSchema),并利用行组的最小值、最大值统计跳过不符合条件的块(条件下推,计划里的 PushedFilters)。再加上按目录拆分的分区,用条件选出来,就可以整个目录地跳过(PartitionFilters)。这三者都只有在优化器能理解条件的时候才会发生。列一旦被套上函数,下推就会消失。 AQE 会在 shuffle 结束之后,根据实际大小修改剩余的计划。所以运行之前的计划(isFinalPlan=false)和运行之后的计划可能不同,真正运行了什么,要在运行之后的计划或事件日志里看。

步骤

  1. 用 /root/spk/plan/convert.py(应用 spk-plan-convert)按 schema 读取订单,加上 order_date(日期)列,以 Parquet 写入 /root/spk/plan/orders_pq。
  2. 在辅助文件 /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()。
  3. 用 /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。
  4. 用 /root/spk/plan/nopush.py(应用 spk-plan-nopush),把同一个问题改成 upper(status) == 'REFUNDED',写入 /root/spk/plan/out/nopush,并把该计划的 PushedFilters 列表写入 /root/spk/plan/out/nopush.json。
  5. 用 /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。
  6. 用 /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。
  7. 用 /root/spk/plan/optimize.py(应用 spk-plan-opt),把串联 qty > 1 和 qty > 3 两个 where、并把 1 + 2 选作 three 的查询的 extended 计划保存到 /root/spk/plan/out/optimized.txt,并执行 collect()。
  8. 在 /root/spk/plan/report.md 中以 ## 밀어내기(韩文,意为“下推”)、## 가지치기(韩文,意为“裁剪”)、## AQE 三节写成报告。第二节放入第 5 步的行数,第三节放入第 6 步的任务数,都用数字写出。

参考

把 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 减少到了几个。