实践了三年多离线数仓把团队主流程从RDD重写成Spark SQL之后我才真正理解“SQL比代码更高效”这句话不是在开玩笑。前两篇写了Spark3.x的核心抽象和数据读写这篇是Spark3.x指北的第三篇专门把Spark SQL讲透它解决了什么问题、一条SQL进到引擎之后经历了什么、3.x版本在SQL层加了哪些关键东西以及我调优和踩坑的真实经验。1. 从RDD到SQLSpark SQL到底解决了什么问题1.1 声明式查询替代命令式计算如果你只用过RDD你的思维模式大概率是“告诉Spark每一步怎么算”先map抽字段再groupByKey聚合最后reduce整理结果。这种命令式写法有几个绕不开的问题第一代码量大一个简单的分组统计要写好几十行第二性能高低完全取决于你个人的优化水平同样的逻辑不同人写出来可能差五倍第三你的每一步计算都被Spark当成“既成事实”优化器没有太多空间帮你做全局调整。Spark SQL的表达方式完全不同。你把需求告诉它剩下怎么执行是引擎的事SELECT city, count(*) AS order_cnt FROM orders GROUP BY city ORDER BY order_cnt DESC这段SQL经过Catalyst优化器之后底层执行计划可能比你手写的RDD版本高效得多——因为优化器能看到整体知道哪些列不需要读、哪些过滤条件可以更早执行、哪些join可以换算法。你给了它“意图”它就还你“最优路径”。1.2 一份代码两种入口DataFrame与SQL的互通很多人以为Spark SQL就是“在Spark里写SQL”其实不止。Spark SQL包含两套入口SQL语言本身以及DataFrame/Dataset编程API。它们共用同一套优化器、同一套执行引擎唯一的区别只是表达形式不同。// 编程API写法 val result spark.read.parquet(/data/orders) .filter($dt 2024-01-01) .groupBy(city) .count() // 等价SQL写法 spark.read.parquet(/data/orders).createOrReplaceTempView(orders) spark.sql( SELECT city, count(*) AS order_cnt FROM orders WHERE dt 2024-01-01 GROUP BY city )实际开发里最常见的做法是行混用复杂的ETL用DataFrame API分析类和报表类需求直接写SQL。Spark SQL的价值不是让你放弃一种写法而是让两种写法共享优化逻辑谁写都行跑起来是一样的执行计划。不管你是从Hive迁移过来的数仓工程师还是在搞实时离线一体化或者刚接触大数据想找个体面的入口Spark SQL都是必须迈过的坎。2. Catalyst优化器一条SQL从文本到执行计划的完整旅程2.1 四条分界线语法解析、语义分析、逻辑优化、物理规划Catalyst是整个Spark SQL的心脏。它做的事情可以类比成编译器你写了一句SQL它把字符串变成一棵语法树再做语义检查、逻辑优化、物理规划最后生成可执行的Java字节码。流程拆分来看就四步SqlParser把SQL文本解析成“未解析的逻辑计划”Unresolved Logical Plan。这个阶段只有语法树不知道表存在不存在、列名对不对。Analyzer带着Catalog去解析列名、表名、函数名完成类型检查和隐式转换输出“解析后的逻辑计划”。你写的dt 2024-01-01在这里会被绑定成某个字段日期类型的比较。Optimizer对解析后的逻辑计划应用一批规则做谓词下推、列剪裁、常量折叠等优化生成“优化的逻辑计划”。SparkPlanner把逻辑计划转换成物理计划再经过Tungsten代码生成变成真正的RDD操作。很多人调优的时候喜欢直接看stage但真正的问题定位一定从物理计划开始。后面第四章我会专门写怎么读执行计划这里先把优化器的规则讲清楚。2.2 几条最常用的优化规则Catalyst优化器内置了大量规则实际工作中让我感觉“值回票价”的是下面这几个谓词下推Predicate PushdownSELECT city, count(*) FROM orders WHERE dt 2024-06-01 GROUP BY city优化器会把WHERE dt 2024-06-01推到数据源扫描阶段。如果底层是Parquet可以跳过不满足条件的RowGroup如果是JDBC数据源这个条件会直接变成SQL的WHERE子句下发给数据库执行。避免把数据全捞回来再过滤。列剪裁Column Pruning列式存储下只读取SQL里真正涉及的列IO量可能减少一个数量级。这个规则对Parquet、ORC这类格式收益极大。你写SELECT city就绝不读amount列。常量折叠Constant FoldingWHERE price 100 50会被写成WHERE price 150省掉每次比较时的重复计算聊胜于无但规则本身展示了优化器的思路任何能在编译期算完的东西绝不留到运行期。Limit下推LIMIT 10会被尽可能推低减少上游处理的数据量。特别是配合排序的场景优化器会识别出“只有TopN”并做相应优化而不是全局排完再取前10。2.3 物理计划与WholeStageCodegen优化器给的是逻辑计划真正执行要靠物理计划。Spark 2.0之后引入的Whole-stage Code Generation会把整个stage内的算子比如Scan、Filter、HashAggregate编译成一段Java代码放进一个方法里执行。这意味着什么在没有代码生成之前Spark的每个算子都是独立的函数调用一行数据在算子间流转需要多次虚函数调用、多次进出内存有了代码生成之后数据在同一个循环体内被连续处理JVM的JIT也能更好地优化。你在执行计划里看到的*(1)前缀代表这个算子属于第1个代码生成阶段*(1) HashAggregate(keys[city#3], functions[count(1)]) - *(1) Project [city#3] - *(1) FileScan parquet orders[city#3,dt#5]代码生成不是银弹但对大部分CPU密集型的聚合、过滤、join场景提升非常明显。了解这个机制之后你再看执行计划就不会只盯着“有没有shuffle”了——算子之间的“融合”程度同样影响性能。3. Spark3.x的SQL层变化ANSI模式、动态分区裁剪与Hive集成3.1 ANSI模式到底改变了什么Spark 3.0开始支持spark.sql.ansi.enabled默认关闭。这个开关直接改变SQL的运行时行为从“容错优先”变成“标准优先”。举几个最常见的差异行为ANSI关闭默认ANSI开启整型溢出返回NULL抛出异常除数为0返回NULL抛出异常类型强转自动尝试转换严格检查失败报错非法日期解析为NULL抛异常关着的时候一条SQL里某个字段溢出只会让那行变NULL任务照跑开着之后遇到脏数据直接让作业失败。看起来ANSI模式“更麻烦”但生产环境里宁可让作业失败也不能让脏数据悄悄混进结果里——尤其是做金融、对账这种场景Null和错误结果比作业失败可怕得多。迁移现有任务到ANSI模式时建议先开着模式跑一遍测试数据把报错逐条修掉。常见问题集中在CAST和日期字符串解析上。3.2 动态分区裁剪星型模型查询的隐形加速器动态分区裁剪Dynamic Partition PruningDPP是Spark 3.0开始成熟的能力默认开启。它的适用场景是数仓里最常见的星型模型一张大事实表一张小维度表通过维度表过滤后join事实表。比如SELECT o.* FROM orders o JOIN dim_city c ON o.city_id c.id WHERE c.region 华东传统静态裁剪只能处理orders表上直接带dt 2024-06-01这种条件。但这里的过滤条件在dim_city表上事实表orders如果按dt分区Spark执行时会在“读取维度表过滤结果”和“扫描事实表”两个stage之间搭建一个动态过滤集扫描事实表时直接跳过不满足条件的分区。我实测过一个业务日订单事实表3亿行按天分区维度表城市表几百行。做完DPP优化之后事实表从扫描30天分区降为扫描10天分区单次任务从22分钟压到9分钟。所以如果你的架构是“大事实表维度表”且发生了join一定要去执行计划里确认DPP有没有生效。有个前提值得一提事实表的连接键和分区键需要是同一套逻辑。比如事实表按dt分区但join条件用的是city_id那DPP也帮不了你只能靠常规的谓词下推和存储裁剪。3.3 Catalog与Hive集成从enableHiveSupport到第三方CatalogSpark SQL从来不只是“自己算一下内存里的表”。真正生产级的用法是跟Hive Metastore打通让Spark能读写Hive数仓表。val spark SparkSession.builder() .appName(spark_sql_practice) .enableHiveSupport() .config(spark.sql.hive.metastore.version, 2.3.9) .config(spark.sql.hive.metastore.jars, maven) .getOrCreate()enableHiveSupport()让Spark接入Hive的Metastore和SerDe体系能建表、能读Hive表、能复用已有的数仓元数据。生产环境还要把hive-site.xml、core-site.xml、hdfs-site.xml放到Spark的classpath下否则连不上MetaStore。如果你部署了Iceberg、Paimon这类湖格式Spark 3.x也给了插件化Catalog的路径。3.4开始支持直接用CREATE CATALOG语法注册外部目录数据湖表可以像普通表一样被Spark SQL查询和写入。这已经不只是数仓场景了而是Spark在做统一的分析入口。4. 靠执行计划调优Spark SQL性能问题的定位路径4.1 学会读EXPLAIN输出定位Spark SQL性能问题第一反应应该是看执行计划而不是盲目调参。spark.sql(SELECT ...).explain(true) // 或者 spark.sql(EXPLAIN EXTENDED SELECT ...)输出里重点看物理计划段几个关键记号记号含义需要警惕FileScan文件扫描如果读了很多列但没有列剪裁检查是不是SELECT *Filter过滤算子过滤出现在扫描之后很远说明谓词下推没生效Exchangeshuffle阶段出现Exchange意味着数据要跨节点传递成本高BroadcastHashJoin广播join小表广播通常很理想SortMergeJoin排序合并join没有广播两边都要shuffle排序大表join的常态WholeStageCodegen*(1)代码生成越多算子融合在一起越好物理执行计划就像一个X光片哪里有问题一目了然全stage就一个FileScan一个大聚合通常没毛病但如果你看到一个小表join大表却走了SortMergeJoin那大概率是广播阈值设置不合适。4.2 三个高频低效模式与计划特征第一个常见问题是扫描层没有裁剪。看到FileScan把表的所有列都读了但后面Project只取了两个字段说明列剪裁没被应用。检查是不是有UDF导致优化器无法推断输出列或者用了SELECT *。第二个常见问题是join没走广播。物理计划里出现SortMergeJoin但其中一边实际数据量很小。这种场景直接调大spark.sql.autoBroadcastJoinThreshold或者写SQL时用/* BROADCAST(t) */强制提示。第三个问题是倾斜导致的单task巨慢。这种在Spark UI的stage详情里表现最明显99%的task早就跑完了只剩一个task跑了一个多小时。执行计划层面看不出太大异常但Exchange后某分区数据量明显不均衡。3.x里用AQE的skew join可以自动拆倾斜分区后面讲参数。4.3 AQE开启后的自动救场与手动兜底Adaptive Query Execution在Spark 3.0引入3.0和3.1默认关闭3.2之后默认开启。它做得最漂亮的一点执行到一半根据实际数据分布动态调整执行计划。参数默认值作用spark.sql.adaptive.enabled3.0/3.1为false3.2起为trueAQE总开关spark.sql.adaptive.coalescePartitions.enabledtrue动态合并shuffle后的分区减少小任务spark.sql.adaptive.advisoryPartitionSizeInBytes64MBshuffle分区目标大小spark.sql.adaptive.autoBroadcastJoinThreshold继承autoBroadcastJoinThreshold动态把SMJ转成BroadcastHJspark.sql.adaptive.skewJoin.enabledtrue自动检测并拆分倾斜join分区这条链路在实际定位时的完整路径我给你还原一下先在Spark UI的SQL页签里找到慢的query点开Physical Plan确认瓶颈是否在某个Exchange、某个Sort、或者某个join用EXPLAIN对比开启AQE前后的物理计划看是否需要手动调参大表join小表还在sort先检查autoBroadcastJoinThreshold和AQE的dynamicJoin再考虑加hint结果末尾还有200个reduce任务在跑但每个都很小打开动态合并shuffle分区调整之后重新跑任务对比stage耗时而不是只看总耗时。AQE还有一个容易忽略的好处它能让写入的表分区数更合理。比如INSERT INTO table SELECT ... GROUP BY keyAQE会按实际数据量帮你合并输出分区减少产生大量小文件的情况。4.4 一个慢查询的完整排查案例有一次我接到一个线上问题每天凌晨的报表任务越跑越慢从最初的15分钟恶化到80分钟。SQL本身不复杂就是三张大表join再group by。我用EXPLAIN看了物理计划发现问题有三个两张事实表join时条件字段类型不一致一个int一个stringCatalyst做了隐式转换导致没有命中bucketing优化所有join都是SortMergeJoin但其中一张子表过滤完之后只有几MB完全可以广播group by stage产生了大量小文件落到HDFS拖慢了下游读数据。当时做了三件事把类型不统一的字段做了CAST统一把spark.sql.autoBroadcastJoinThreshold从10MB调到100MB让子表走BroadcastHashJoin给写入侧加了repartition分区控制避免再产生几百个小文件。调整之后任务从80分钟恢复到12分钟左右。这个案例里没有任何一个“高深”动作就是从执行计划里看到问题然后做出针对性调整。5. 日常开发里的几个高频坑与绕坑姿势5.1 小文件问题Spark SQL也会“制造垃圾”写完INSERT INTO table SELECT ...之后最怕看到的结果就是表目录下多出来200个几KB的文件。原因通常是spark.sql.shuffle.partitions保持默认200而写入模式是shuffle。解决路径有两种写coalesce或repartition控制并行度或者依赖AQE的coalescePartitions自动合并小分区。更规范的做法是设计阶段就用分桶表bucket table让相同key的数据落进固定数量的文件。判断标准我一般这样控制单文件在64MB到128MB之间比较理想。文件太小的直接后果是下游扫描时要起一堆task元数据压力也大。5.2 UDF序列化、确定性与性能三重坑写自定义UDF是Spark SQL开发里最常见的坑源。第一坑是序列化你在UDF闭包里引用了无法序列化的对象跑起来直接NotSerializableException。第二个坑是确定性如果UDF内部依赖随机数或外部状态优化器会放弃对它的下推、缓存等很多优化。第三个坑是性能JVM内部的UDF是逐行调用的一个简单的Java UDF跑一亿行函数调用开销会大得惊人。能解决先解决优先用内置函数实在要写3.x里优先考虑Pandas UDF向量化执行而不是普通UDF。我做过一次对比同一个字符串处理的UDFPandas UDF比逐行Java UDF快了接近一个数量级。5.3 NULL与类型隐式转换掉数据常常因为你没看见Spark SQL里NULL的语义很容易让人踩坑。WHERE column NULL永远为false必须用IS NULL。两列join时如果一边有NULL这些行会被直接丢掉还不报错。很多“数据少了几万行”的血案查到最后都是join key里有NULL。类型隐式转换更隐蔽。int和string在join时Spark会尝试转成公共类型但转换会破坏分区裁剪、bucketing优化甚至触发整个表扫描。写join条件时老老实实把两边类型统一成相同类型收益非常直接。5.4 缓存用对了是加速器用错了是内存杀手df.cache()不是万能的。缓存适合在同一个SparkSession里被多个查询复用、且表数据本身不大的场景。如果一张表只是跑一次性查询缓存反而浪费内存和序列化时间。缓存之后记得看StorageLevel默认MEMORY_AND_DISK会优先放内存内存不够落磁盘生产环境可以接受但如果缓存大表且executor内存本就不宽裕可能挤掉其他task的堆内存引发频繁GC甚至OOM。一个个人习惯缓存的表尽量先filter和select之后再cache缓存小结果集而不是大原表这比调任何参数都见效。5.5 写在最后的实操体会如果只挑一条经验分享那就是Spark SQL的一切问题都要回到执行计划上去回答。慢就先问是不是裁剪没生效数据不对就先看join key有没有NULL和类型不一致小文件多就先想shuffle分区和写入方式。把这些习惯养成了你对Spark SQL的掌控力会上一个台阶。Spark3.x的AQE和DPP已经帮我们挡掉了不少原本要手动做的优化但工具帮不了的地方还是要靠你对执行引擎的理解来兜底。