Spark Catalyst优化器如何选择物理计划?代价函数是什么?
Spark Catalyst优化器的执行流程与代价优化详解
一、Catalyst的核心优化执行流程
Catalyst的优化是分阶段的树形转换过程,基于抽象语法树(AST)处理,核心分为四个阶段:
- 解析(Parsing):将SQL或DataSet代码转换为未解析的逻辑计划,仅做语法校验,不验证表、列等元数据是否存在。
- 绑定(Analysis):结合元数据目录(Catalog),把未解析计划转换为已解析的逻辑计划,完成列存在性、数据类型匹配等校验。
- 逻辑优化(Logical Optimization):基于规则(RBO)对逻辑计划进行批量转换,比如谓词下推、列裁剪、常量折叠、聚合下推等,反复迭代直到无优化空间。
- 物理计划(Physical Planning):将优化后的逻辑计划转换为多个可选物理执行方案,再通过代价优化(CBO)选出最优计划。
二、物理计划生成与最优选择逻辑
1. 物理计划生成逻辑
逻辑计划是抽象计算逻辑,物理计划则是具体执行策略。Catalyst会为每个逻辑算子生成多种可行的物理实现:
- 针对Join算子,生成Broadcast Hash Join、Shuffle Hash Join、Sort Merge Join等方案;
- 针对聚合算子,生成HashAggregate、SortAggregate两种实现。
生成过程基于规则映射:每个逻辑算子对应一组可转换的物理算子,遍历所有组合后生成完整的物理计划树。
2. 代价优化(CBO)核心逻辑
Catalyst的CBO依赖统计信息计算物理计划代价,最终选择代价最低的方案:
统计信息来源
- 表级统计:总行数、数据总大小、平均行宽;
- 列级统计:列基数(不同值数量)、空值占比、极值、直方图(Spark 2.3+支持);
需通过ANALYZE TABLE [table_name] COMPUTE STATISTICS或ANALYZE TABLE [table_name] COMPUTE STATISTICS FOR COLUMNS [col1, col2]提前生成,存储在Catalog中。
代价计算维度
主要从两个维度计算:
- I/O代价:数据读取、shuffle的磁盘/网络IO量,比如Join时shuffle的数据规模;
- CPU代价:算子计算开销,比如排序、哈希匹配的CPU消耗。
每个物理算子都有对应代价公式,比如Sort Merge Join的代价包含排序CPU开销+shuffle IO开销,Broadcast Hash Join则包含广播数据网络开销+哈希匹配CPU开销。
最优计划选择
Catalyst遍历所有物理计划,计算每个计划的总代价(I/O与CPU的加权和),选择总代价最小的作为最终执行计划。若无统计信息,会退化为RBO模式,使用默认物理算子(如默认Sort Merge Join)。
三、典型代价优化场景
- 小表与大表Join时,CBO计算小表数据量,判断是否适合广播,选择Broadcast Hash Join避免大表shuffle;
- 聚合列基数较低时,选择HashAggregate而非SortAggregate,降低CPU开销;
- 过滤条件选择性高时,将谓词下推至数据源端(如Parquet/Orc),减少读取数据量。
内容的提问来源于stack exchange,提问作者Hijaw
相关产品推荐
相关产品推荐

