You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark相似逻辑与物理查询计划的优劣评估及选择咨询

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)

我的理解与疑问

我认为优化方案更优,理由为:

  1. Filter执行时机更早(优化计划步骤2 vs 原计划步骤8)
  2. 提前过滤数据,减少后续pivot等操作的数据量

请问我的理解是否正确?(注:已知AdaptiveSparkPlan会在运行时调整执行计划,目前暂聚焦于当前计划的评估)


评估结论

你的理解完全正确,优化后的查询计划确实更高效,核心原因如下:

  1. 过滤前置大幅降低无效计算
    原计划中Filter在聚合(SortAggregate)之后执行,意味着先对全量数据完成了pivot所需的分组、排序、聚合操作,再过滤掉不需要的行——这会浪费大量计算资源在无效数据上。而优化计划将Filter放在Scan之后、Project之前,直接在读取数据阶段就过滤掉col2不属于['value1','value2']的行,从源头缩减了后续所有处理步骤的数据量。

  2. 数据量缩减的连锁优化效应
    虽然统计数据中Filter后的sizeInBytes显示仍为3.2GiB(大概率是因为缺少列统计信息导致估算不准),但实际执行中,提前过滤会减少后续Project、Sort、Exchange、Aggregate等步骤处理的数据行数,降低Shuffle阶段的网络传输开销,同时减少排序和聚合的计算时间。

  3. 代码简化减少执行计划冗余
    优化后的代码去掉了冗余的withColumn和dropna操作,直接通过where提前过滤,再用selectExpr重命名列,逻辑更简洁,也避免了原计划中dropna对应的后续Filter操作,减少了执行计划的冗余步骤。

补充说明:即便AdaptiveSparkPlan会在运行时调整执行计划,当前的初始计划优化也为运行时调整打下了更好的基础——初始处理的数据量越小,ASAP能做的动态优化(比如调整Shuffle分区、合并小任务)效率也会越高。


内容的提问来源于stack exchange,提问作者Alex

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 19:04:54