PySpark相似逻辑与物理查询计划的优劣评估及选择咨询
问题背景
我需要评估两个输出结果一致的PySpark查询计划,此前自行调研未找到有效信息。两个计划的核心逻辑一致:
- 选择特定列
- 对特定值执行pivot操作并计算聚合
- 基于pivot结果生成新列
原代码与物理执行计划
原代码
attributes = spark.read.format('delta').load(path).select('col1', 'col2', 'col3') attributes = attributes.groupBy('col1').pivot('col2', values = ['value1', 'value2']).agg(max('col3')) attributes = attributes.withColumn('new', col('pivotvalue')).withColumn('new2', col('pivotvalue2')) attributes = attributes.select('col1', 'new', 'new2') attributes = attributes.dropna(how = 'all', subset = ['new', 'new2']) # The drop Na is because the values that are not in the set given to the pivot column are not required.
原物理执行计划
== Physical Plan == AdaptiveSparkPlan (10) +- == Initial Plan == Project (9) +- Filter (8) +- SortAggregate (7) +- Sort (6) +- Exchange (5) +- SortAggregate (4) +- Sort (3) +- Project (2) +- Scan parquet (1)
优化后代码与物理执行计划
优化后代码
attributes = spark.read.format('delta').load(path).select('col1', 'col2', 'col3').where("col2 in ('value1', 'value2')") attributes = attributes.groupBy('col1').pivot('col2').agg(max('col3')).selectExpr('col1', 'new as pivotvalue', 'new2 as pivotvalue2')
优化后物理执行计划
== Physical Plan == AdaptiveSparkPlan (10) +- == Initial Plan == Project (9) +- SortAggregate (8) +- Sort (7) +- Exchange (6) +- SortAggregate (5) +- Sort (4) +- Project (3) +- Filter (2) +- Scan parquet (1)
成本统计数据对比
原计划统计
1. Relation -> Statistics(sizeInBytes=3.2 GiB, ColumnStat: N/A) 2. Project -> Statistics(sizeInBytes=430.2 MiB, ColumnStat: N/A) 3. Aggregate -> Statistics(sizeInBytes=1147.1 MiB, ColumnStat: N/A) 4. Filter -> Statistics(sizeInBytes=1147.1 MiB, ColumnStat: N/A) 5. Project -> Statistics(sizeInBytes=645.2 MiB, ColumnStat: N/A)
优化计划统计
1. Relation -> Statistics(sizeInBytes=3.2 GiB, ColumnStat: N/A) 2. Filter -> Statistics(sizeInBytes=3.2 GiB, ColumnStat: N/A) 3. Project -> Statistics(sizeInBytes=430.2 MiB, ColumnStat: N/A) 4. Aggregate -> Statistics(sizeInBytes=1147.1 MiB, ColumnStat: N/A) 5. Project -> Statistics(sizeInBytes=645.2 MiB, ColumnStat: N/A)
我的理解与疑问
我认为优化方案更优,理由为:
- Filter执行时机更早(优化计划步骤2 vs 原计划步骤8)
- 提前过滤数据,减少后续pivot等操作的数据量
请问我的理解是否正确?(注:已知AdaptiveSparkPlan会在运行时调整执行计划,目前暂聚焦于当前计划的评估)
评估结论
你的理解完全正确,优化后的查询计划确实更高效,核心原因如下:
过滤前置大幅降低无效计算
原计划中Filter在聚合(SortAggregate)之后执行,意味着先对全量数据完成了pivot所需的分组、排序、聚合操作,再过滤掉不需要的行——这会浪费大量计算资源在无效数据上。而优化计划将Filter放在Scan之后、Project之前,直接在读取数据阶段就过滤掉col2不属于['value1','value2']的行,从源头缩减了后续所有处理步骤的数据量。数据量缩减的连锁优化效应
虽然统计数据中Filter后的sizeInBytes显示仍为3.2GiB(大概率是因为缺少列统计信息导致估算不准),但实际执行中,提前过滤会减少后续Project、Sort、Exchange、Aggregate等步骤处理的数据行数,降低Shuffle阶段的网络传输开销,同时减少排序和聚合的计算时间。代码简化减少执行计划冗余
优化后的代码去掉了冗余的withColumn和dropna操作,直接通过where提前过滤,再用selectExpr重命名列,逻辑更简洁,也避免了原计划中dropna对应的后续Filter操作,减少了执行计划的冗余步骤。
补充说明:即便AdaptiveSparkPlan会在运行时调整执行计划,当前的初始计划优化也为运行时调整打下了更好的基础——初始处理的数据量越小,ASAP能做的动态优化(比如调整Shuffle分区、合并小任务)效率也会越高。
内容的提问来源于stack exchange,提问作者Alex

