1. 先说清楚:SparkSQL到底替你扛了哪些脏活累活
下午刚帮同事排查完一个跑了40分钟的Spark任务,最后定位到问题出在一段写得极其绕的DataFrame API链式调用上——我把它改写成三条SQL,执行时间掉到7分钟。这种场景在我手里已经发生过无数次,也让我越来越确信一件事:在大数据开发这个岗位上,SparkSQL不只是一门"要学的技术",它是生产环境里最能直接帮你省时间的工具。
这篇东西我不会给你贴官方文档式的API清单,那种东西你上网一搜一大把。我打算从实际干活的角度,把SparkSQL里那些高频操作掰开揉碎了讲一遍:哪些写法在集群上会踩坑、哪些优化手段是真有用的、哪些坑我已经替你踩过了。适合正在做离线数仓、实时数仓或者ETL开发的朋友,尤其是那种刚把Spark跑起来、天天被任务性能和复杂逻辑折磨的阶段。不管你是写Scala还是PySpark,核心思路都一样,代码我尽量两种都给。
先立个flag:看完这篇,你至少应该能完成"读原始日志 -> 清洗 -> 关联维表 -> 聚合 -> 落结果表"这条完整链路,而且知道自己写出的每个操作大概会跑成什么样。
2. 动手前的第一件事:把SparkSession和基础读写彻底搞明白
很多人上来就写spark.sql("select ..."),结果第一步就卡住——spark这个对象从哪来的?在spark-shell里它是现成的,但你写独立作业的时候必须自己构建。这一步看着简单,里面埋着三个我见过无数次的低级错误。
2.1 构建SparkSession的推荐姿势
val spark = SparkSession.builder() .appName("example_etl") .config("spark.sql.shuffle.partitions", "200") .config("spark.sql.adaptive.enabled", "true") .enableHiveSupport() // 如果要读写Hive表,这行必须有 .getOrCreate()enableHiveSupport()这行我单独提出来说。我碰到过不止一个项目,代码里不写这行,然后跑到spark.sql("select * from ods.table")的时候就报Table or view not found。原因很简单:没有启用Hive支持,SparkSQL就不知道去哪找Hive的元数据。如果你只是读写HDFS上的文件或者临时表,那无所谓;但只要你的数据是正经落在Hive表里,这一行必须加。
PySpark版本长这样:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("example_etl") \ .config("spark.sql.shuffle.partitions", "200") \ .config("spark.sql.adaptive.enabled", "true") \ .enableHiveSupport() \ .getOrCreate()2.2 读文件阶段最容易犯的错:Schema推断
读CSV就开始写:
val df = spark.read .option("header", "true") .csv("hdfs:///data/raw/20240101/")这个写法在本地小文件上没问题,一上集群参与大规模处理就会出现两个毛病。第一,Spark需要扫一遍全部数据才能推断出每列的类型,数据量一大,光Schema推断就能给你多耗出十几分钟。第二,推断出来的类型大概率是错的——比如某天数据里金额字段全是整数,它推断成LongType,第二天数据里多出几条带小数的,任务直接报错。
更稳的姿势是显式指定Schema:
import org.apache.spark.sql.types._ val schema = StructType(Array( StructField("order_id", StringType, true), StructField("user_id", LongType, true), StructField("amount", DoubleType, true), StructField("ts", TimestampType, true) )) val df = spark.read .option("header", "true") .schema(schema) .csv("hdfs:///data/raw/20240101/")显式指定Schema的核心好处是:任务跑起来之前,执行计划就已经确定了,不需要额外的扫描来推断类型;而且数据类型可控,脏数据会在转换阶段就暴露,不会跑到半路才炸。
顺带提一句Parquet。如果你有选择权,结果数据尽量用Parquet落盘。列式存储+压缩,查询效率比CSV高出一截,而且Parquet自带Schema,读的时候不需要指定,Spark直接能从文件元数据里拿到字段结构。这个习惯我在团队里强推了很久,谁用谁知道。
2.3 临时视图:SQL和DataFrame API之间的桥梁
我刚入行的时候特别喜欢纯写DataFrame API,觉得类型安全、IDE有提示,后来发现一个尴尬的事实:项目里很多人不会DataFrame API,只会写SQL。而且坦白说,有些复杂逻辑用SQL表达比用API链式调用清晰得多。解决方案就是临时视图。
df.createOrReplaceTempView("tmp_order") spark.sql(""" SELECT user_id, SUM(amount) AS total_amount FROM tmp_order WHERE dt = '20240101' GROUP BY user_id """).show()注意区分createOrReplaceTempView和createGlobalTempView。前者是session级别的,出了这个SparkSession就没了;后者是跨session共享的,访问的时候要带前缀global_temp。我们生产上基本只用session级别的,全局视图用到的场景极少,还容易引起命名冲突。
3. 查数、过滤、聚合、关联:日常用得最狠的几板斧
这一节我按使用频率排序来讲。先说结论:一个正常数仓开发每天写的SparkSQL,80%逃不出SELECT + WHERE + GROUP BY + JOIN。但就是这几板斧,细节上的差距决定了你的任务是跑5分钟还是跑50分钟。
3.1 过滤条件里藏着性能黑洞
从一个真实案例说起。有个同事处理2亿行订单数据,代码是这样的:
val filtered = df.filter("dt = '20240101' AND amount > 100")跑得很慢,看执行计划发现Scan阶段扫了全表所有分区。原因特别简单:这张表是按dt做分区列的,Hive表里天然能通过分区裁剪只读当天数据,但Spark读的是原始文件路径,不知道这个奇偶关系。如果dt是分区字段,你必须确保过滤条件下推:
SELECT * FROM table WHERE dt = '20240101'这条SQL之所以走分区裁剪,是因为Spark的Catalyst优化器能够识别dt是分区列,然后触达执行计划阶段就只读取对应分区的文件。但如果你在DataFrame API里先把整个DataFrame读进来再filter,那文件扫描这一步已经发生了,优化器再厉害也救不回来。
区分两个概念:分区裁剪(只读相关分区文件)和谓词下推(把过滤条件推到读取阶段,减少进入内存的数据量)。两个都做到,才算把过滤写明白了。
3.2 JOIN的三种常见姿势和实务选择
SparkSQL里的JOIN,大方向分三种:
| JOIN类型 | 触发条件 | 特点 |
|---|---|---|
| Broadcast Join | 小表小于 spark.sql.autoBroadcastJoinThreshold(默认10MB) | 把小表广播到每个Executor,不走Shuffle,极快 |
| Sort Merge Join | 两表都较大且无广播条件 | 按key排序后合并,有Shuffle,慢但稳 |
| Shuffle Hash Join | 关闭SortMerge或特定参数下 | 按key哈希分桶后建哈希表,内存压力大 |
实战里最常见的优化就是把参与Join的小表size控制住,让它走Broadcast Join。比如关联维表:
val dimDF = spark.read.parquet("hdfs:///dim/user_info/") val orderDF = spark.read.parquet("hdfs:///ods/order/") val joined = orderDF.join(dimDF, Seq("user_id"), "left")如果dimDF只有几万条,优化器大概率自动广播。但如果维表有200MB,默认阈值10MB不满足,启动计划就会变成SortMergeJoin。这时候你可以显式提示:
val joined = orderDF.join(broadcast(dimDF), Seq("user_id"), "left")broadcast函数显式标记小表,告诉优化器"我确认这张表适合广播,别犹豫了"。不过要小心,如果这张表实际很大还要强制广播,Executor内存直接爆掉,OOM之后任务反复restart,这个坑我踩过一次后就再也不敢乱标了。准确的做法是先看表的实际大小,再决定是否广播。
3.3 GROUP BY聚合:别小看group by后边的坑
聚合操作本身没太多戏,但聚合之前的数据倾斜问题太常见了。举个具体场景:
SELECT province, COUNT(*) AS cnt FROM order_info GROUP BY province如果某个省份的订单量占了一半,就会有一个Reduce任务处理的数据量远超其他任务,这一轮跑得慢,整个Stage都被拖住。这就是最典型的数据倾斜。
针对这种“热点key”倾斜,实战里有一个挺好用的思路:两阶段聚合。
-- 第一步:先给key加随机前缀,打散热点 SELECT province_pre, COUNT(*) AS cnt_pre FROM ( SELECT CASE WHEN province = 'ZHEJIANG' THEN CONCAT('ZHEJIANG_', FLOOR(RAND() * 10)) ELSE province END AS province_pre FROM order_info ) t GROUP BY province_pre -- 第二步:去掉前缀,重新聚合 SELECT CASE WHEN province_pre LIKE 'ZHEJIANG_%' THEN 'ZHEJIANG' ELSE province_pre END AS province, SUM(cnt_pre) AS cnt FROM tmp_result GROUP BY province思路是把倾斜的key先加随机后缀拆成10个分桶,让它们分散到不同Reduce任务里,第一阶段先各自算,第二阶段再汇总。这个法子不解决所有倾斜场景,但对高频key倾斜非常有效。
3.4 窗口函数:排序、去重、取TopN的利器
窗口函数在生产里用得太频繁了,我说三个高频场景。
场景一:同key内按时间取最新一条。
SELECT * FROM ( SELECT *, ROW_NUMBER() OVER(PARTITION BY order_id ORDER BY update_time DESC) AS rn FROM order_log ) t WHERE rn = 1这个写法是做数据去重、保留最新状态的通用解法,比如订单更新日志按order_id去重取最新状态。注意ROW_NUMBER()的PARTITION BY和ORDER BY顺序不要搞反,先定分组再定组内排序。
场景二:分组TopN。
SELECT user_id, category, amount FROM ( SELECT user_id, category, amount, RANK() OVER(PARTITION BY user_id ORDER BY amount DESC) AS rk FROM order_info ) t WHERE rk <= 3RANK()和DENSE_RANK()的区别在于并列名次是否占用后续名词,业务上"取前三"建议想清楚用哪个。
场景三:累加趋势。
SELECT dt, pay_amount, SUM(pay_amount) OVER(ORDER BY dt) AS cumulative_amount FROM daily_pay这种写法在做GMV累计、用户生命周期价值分析的时候非常常用。窗口函数的好处是代码简洁,但要注意:如果窗口的ORDER BY或PARTITION BY字段分布不均匀,同样会产生数据倾斜。性能优化和功能实现得同时考量。
3.5 UDF:能不用就尽量别用,但用的时候要写对
很多业务逻辑用纯SQL表达不了,比如JSON解析、IP地址转地理位置、复杂字符串处理。这时候就需要UDF。
import org.apache.spark.sql.functions.udf val parseRegion = udf((ip: String) => { // 实际逻辑省略,这里做IP解析 "province" }) val df2 = df.withColumn("region", parseRegion(col("client_ip")))UDF的坑主要在性能上:普通UDF每行数据都要经过JVM和Spark内部之间的序列化/反序列化,行数一多开销非常明显。如果逻辑不复杂,优先考虑用内置函数替代。比如JSON解析可以用get_json_object,字符串截取可以用substring,这些内置函数是Catalyst优化器能直接处理的,性能比UDF高一个量级。
如果确实非用UDF不可,写的时候注意:不要在UDF内部创建单例对象以外的重量级资源,比如每个函数调用都连一次数据库,这种写法跑2亿行就会连2亿次数据库,任务永远跑不完。
4. 为什么有时"烂SQL"跑得也快:Catalyst优化器到底替你做了什么事
这个部分有点偏原理,但我觉得搞懂一点执行计划的生成逻辑,对写SparkSQL有质的提升。不然你永远在"试错":同一个结果,换一种写法,性能差异巨大,你只能靠猜。
4.1 从SQL到执行计划,中间隔着一套"逻辑改写"
SparkSQL跑一条查询时,经历的过程大致是:
SQL语句 -> 语法解析 -> 生成未优化的逻辑计划 -> Catalyst优化器做规则优化 -> 生成物理计划 -> 转成RDD作业执行。
Catalyst优化器会做几件特别关键的事:
- 谓词下推:把
WHERE条件尽可能提前到读取数据的时候执行,减少进入计算引擎的数据量。 - 列剪枝:只读查询里用到的列,不用的列直接跳过,对列式存储格式(Parquet)尤其有效。
- 常量替换:
WHERE dt = '20240101'里的字符串常量会在优化阶段被推导进分区裁剪。 - Join重排:多个表关联时,优化器会尝试把小表先关联,减少中间数据量。
这就是为什么有时候你在DataFrame API里写了一长串filter().join().groupBy(),它跑起来依然不慢——因为优化器把你那串链式调用转换成的逻辑计划,经过了这么多轮规则改写,最后变成的物理执行可能和你写的顺序完全不同。
4.2 但优化器不是万能的,这三个"盲区"你得自己补
盲区一:你自己写的笛卡尔积,优化器不敢乱动。cross join就是横着连,优化器没那个胆量自动帮你加关联条件,因为它不知道你的业务意图。所以写JOIN一定要写ON条件,不要只写WHERE。
盲区二:UDF内部的逻辑优化器完全看不到。前面说了,UDF对优化器是个黑盒,它不知道你UDF里能不能下推、能不能剪枝,所有优化都失效。能用内置函数绝不用UDF,一部分原因就在这里。
盲区三:三张以上表JOIN的时候,关联顺序优化依赖统计信息。如果表没有收集过统计信息,优化器做Cost-Based Optimization时没有参考数据,它只能基于经验猜测,猜错了执行计划就烂了。生产环境里定期跑ANALYZE TABLE更新统计信息,不是一个可选项,是必要的维护工作。
4.3 学会用EXPLAIN看执行计划,别再瞎猜性能瓶颈
定位性能问题最快的方式是看执行计划,而不是凭经验乱猜。SparkSQL提供了EXPLAIN命令。
EXPLAIN SELECT user_id, SUM(amount) FROM tmp_order WHERE dt = '20240101' GROUP BY user_id实际执行时可以看带物理计划细节的完整输出,注意几个关键节点:
Scan阶段的Partition Count:理想情况下应该是被裁剪后的分区数,如果这里还写着全部分区,说明分区裁剪没生效。Exchange节点:出现这个意味着有Shuffle,Shuffle量越大,任务越慢。HashAggregatevsSortAggregate:前者比后者效率高,如果执行计划里出现了SortAggregate,可以考虑开spark.sql.legacy.allowHashOnMapType之类的参数,或者给聚合字段加合理排序。
我一直跟组里的新人说一句话:先学会读执行计划,再谈性能优化。不然你连问题出在哪都不知道,优化就是撞大运。
5. 从30分钟压到6分钟:一次完整的生产任务性能排查实录
前面讲了原理,这节我拿一个真实的排查案例来串一遍。完整地看一遍定位思路,比记住一百个孤立参数有用。
5.1 初始状态:任务跑了30分钟,卡在一个奇怪的Stage
当时那个任务逻辑不复杂:一天订单明细数据,大概3亿行,关联一张8000万行的用户维表,聚合后输出到一个Hive分区表。跑了30分钟出头,看Spark UI,发现大部分时间耗在第5个Stage,Stage里所有Task的输入端数据量差异巨大,最大的Task处理了接近80%的数据。
我当时的第一判断:SortMergeJoin阶段发生了严重的数据倾斜。倾斜的来源基本可以确定是关联字段user_id的分布不均匀——少量高频用户贡献了大量的订单。
5.2 排查链路:从Spark UI到执行计划再到代码
第一步,看Spark UI里Stage 5的详情。输入数据规模那一栏,几个Task的输入量是几百GB对几十MB的差距,倾斜实锤。
第二步,找到倾斜发生在哪个算子。Stage 5对应的算子就是Exchange(Shuffle)。Shuffle的key是user_id,自然就是用户维表关联的时候,热key全部进了某一个Reduce的桶。
第三步,回看代码。关联条件是order.order_id = user.id(注意这里写错了,应该是order.user_id = user.id,是同事笔误,但这类看似天真的错误在真实代码里就是会出现),不过这个错误先不提。核心问题是大量订单带过来的时候,同一个user_id的订单全压在reduce一侧。
第四步,选择方案。当时维表8000万行,不能直接广播;但可以试试分桶广播的思路,把高频key拆开。实际选择的方案是加盐两阶段聚合:给关联key加上随机前缀,打散之后先做一次不完全的关联和聚合,再按真实key做最终聚合。配合spark.sql.shuffle.partitions从默认200调高到600,让每个Task处理的数据量进一步下降。
5.3 改动后的变化:为什么从3亿行翻倍到6亿行,时间反而缩短了
改写后,中间数据量从3亿变成了大约6亿(因为加了盐之后数据膨胀了),但Shuffle分布均匀了,原本卡在最慢Task上的时间瓶颈没有了。整个任务从30分钟降到了9分钟。后来又进一步做了三件事:
- 把维表在关联前用
repartition按user_id做了哈希预分区,减少Sort合并时的数据移动。 - 打开了
spark.sql.adaptive.coalescePartitions.enabled,让AQE自动合并小分区,避免最后输出的Reducer太多造成小文件问题。 - 给维表加了一层缓存,虽然不是决定性因素,但第二天的重跑任务快了约15%。
最终任务稳定在6分钟出头。
5.4 这件事反映出的几个通识
第一,数据倾斜不是玄学,是可以定位的。Spark UI的Stage详情就是你的第一手现场证据,别一上来就改代码。第二,改写方案要先想清楚中间数据量的变化。加盐聚合一定会让中间数据变多,但只要能换来均匀分布,总时长大概率是降的。第三,AQE(自适应查询优化)能帮你兜底,Reduce端的分区合并和倾斜自动处理都是好东西,建议直接把spark.sql.adaptive.enabled设为true,现在的Spark 3.x默认是开的,别手贱关掉。
6. 从一个原始日志文件到一张主题宽表:完整链路实操
理论部分差不多了,最后拿一条完整的ETL链路把前面讲的内容串起来。这个场景很有代表性:原始数据是JSON格式的行为日志,需要清洗、解析、关联维表、聚合,最后落到一张按天分区的Hive表。
6.1 场景定义:行为日志到用户主题宽表
假设原始日志长这样:
{"user_id": 12345, "action": "click", "item_id": "sku_9981", "ts": "2024-01-01 12:23:45", "extra_info": "{\"source\": \"homepage\", \"ab_test\": \"groupA\"}"}目标表结构(用户主题宽表):user_id, total_click, total_buy, last_active_date, favorite_category。
6.2 第一步:读取和解析
val raw = spark.read.textFile("hdfs:///data/log/20240101/")先用textFile按行读,因为这种非结构化JSON用spark.read.json()直接读,Schema推断很不可控。读进来后在SQL里解析:
CREATE OR REPLACE TEMP VIEW parsed AS SELECT get_json_object(value, '$.user_id') AS user_id, get_json_object(value, '$.action') AS action, get_json_object(value, '$.item_id') AS item_id, get_json_object(value, '$.ts') AS ts, get_json_object(value, '$.extra_info.source') AS source FROM raw_logget_json_object这个内置函数处理JSON字段非常稳,直接用它而不是写UDF,性能差距明显。
6.3 第二步:按用户聚合
CREATE OR REPLACE TEMP VIEW user_agg AS SELECT user_id, SUM(CASE WHEN action = 'click' THEN 1 ELSE 0 END) AS total_click, SUM(CASE WHEN action = 'buy' THEN 1 ELSE 0 END) AS total_buy, MAX(ts) AS last_active_date FROM parsed GROUP BY user_id这里用了SUM(CASE WHEN)做条件计数,比COUNT(FILTER...)的写法更加通用,而且SparkSQL对这两种写法的优化基本等价,选哪种纯看个人口味。
6.4 第三步:关联维表取偏好品类
CREATE OR REPLACE TEMP VIEW user_final AS SELECT a.user_id, a.total_click, a.total_buy, a.last_active_date, b.favorite_category FROM user_agg a LEFT JOIN dim_user_behavior b ON a.user_id = b.user_id如果dim_user_behavior表不大,记得用broadcast;如果太大,至少要确保它按user_id做了合理的文件组织,避免全表扫描。
6.5 第四步:写结果表,注意动态分区
最后落Hive表,假设目标表已经有分区字段dt:
spark.sql(""" INSERT OVERWRITE TABLE dws.user_wide PARTITION(dt='20240101') SELECT user_id, total_click, total_buy, last_active_date, favorite_category FROM user_final """)这里要提醒一下动态分区的坑。如果你要按多个字段动态分区,得先打开参数:
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")默认INSERT OVERWRITE静态分区会先删掉整个分区再写入,如果业务上有"只想覆盖某几个子分区"的需求,必须开动态分区模式。这个参数不打开,你覆盖写的时候会把整个目标分区清掉,数据出事故了我见得太多了。
6.6 这条链路里容易忽略的两个细节
小文件问题。聚合结果如果分区数太多,落盘时会生成大量小文件,后续查询会被拖死。解决办法是最后写结果前做一次repartition或coalesce控制输出分区数,比如:
df.repartition(50).write.mode("overwrite").insertInto("dws.user_wide")输出格式。如果目标表是Hive外部表且数据量很大,考虑用Parquet + Snappy压缩,比文本格式减少约70%的存储空间,查询速度也更快。
7. 需要特别注意的另一个方向:SparkSQL写法的可读性和可维护性
写完功能、调完性能之后,还有一个经常被忽略的问题:这段代码三个月后别人能不能看懂。我在团队里审核代码的时候,看到过太多那种功能正确但完全读不动的SparkSQL。
几个建议,都是自己在生产环境里吃过亏总结的:
给临时视图起有业务含义的名字。tmp1、tmp2这种临时视图,一多起来就是灾难。哪怕多打几个字,起成daily_order_agg、user_base_info,后面排查问题时省的事远比你敲字的时间值钱。
复杂逻辑拆成多段中间视图,不要一个超级SQL搞定一切。这个和写普通SQL的习惯一样,一个SQL动辄两三百行,出了问题根本没法定位。拆成raw_parsed -> cleaned -> enriched -> aggregated,每层都能单独验证。
该注释的地方一定要注释。特别是那种优化型的写法,比如加盐聚合、广播小表,如果不在代码里说明"为什么这么写",后人看不懂可能直接给你改成普通写法,性能一夜回到解放前。
结果是数字的字段,类型要统一。SparkSQL对Long和Double混着用会出现精度问题,特别是金额相关的字段,前期类型定义不统一,后面聚合结果对不上账,排查两个小时都为这个。
这些都是"看起来不影响功能,实际影响很大"的事。写SparkSQL不是写给机器看的,是写给下一个维护者看的,包括三个月后的你自己。
文章写到这里,技术内容基本都覆盖了。我最后想说的是:SparkSQL的上手门槛其实不高,SQL语法你本来就会,真正的深度在"懂原理"和"会排查"这两件事上。建议你拿到一个任务后,别急着写完代码就跑,先想清楚:数据量级多大、关联的表多大、哪里可能倾斜、分区裁剪能不能生效。习惯养成了,你的任务会跑得比别人稳很多。