ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

SparkSQL性能优化实战:从数据倾斜到Catalyst优化器关键技巧

SparkSQL性能优化实战:从数据倾斜到Catalyst优化器关键技巧 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第二天数据里多出几条带小数的任务直接报错。更稳的姿势是显式指定Schemaimport 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. 查数、过滤、聚合、关联日常用得最狠的几板斧这一节我按使用频率排序来讲。先说结论一个正常数仓开发每天写的SparkSQL80%逃不出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节点出现这个意味着有ShuffleShuffle量越大任务越慢。HashAggregatevsSortAggregate前者比后者效率高如果执行计划里出现了SortAggregate可以考虑开spark.sql.legacy.allowHashOnMapType之类的参数或者给聚合字段加合理排序。我一直跟组里的新人说一句话先学会读执行计划再谈性能优化。不然你连问题出在哪都不知道优化就是撞大运。5. 从30分钟压到6分钟一次完整的生产任务性能排查实录前面讲了原理这节我拿一个真实的排查案例来串一遍。完整地看一遍定位思路比记住一百个孤立参数有用。5.1 初始状态任务跑了30分钟卡在一个奇怪的Stage当时那个任务逻辑不复杂一天订单明细数据大概3亿行关联一张8000万行的用户维表聚合后输出到一个Hive分区表。跑了30分钟出头看Spark UI发现大部分时间耗在第5个StageStage里所有Task的输入端数据量差异巨大最大的Task处理了接近80%的数据。我当时的第一判断SortMergeJoin阶段发生了严重的数据倾斜。倾斜的来源基本可以确定是关联字段user_id的分布不均匀——少量高频用户贡献了大量的订单。5.2 排查链路从Spark UI到执行计划再到代码第一步看Spark UI里Stage 5的详情。输入数据规模那一栏几个Task的输入量是几百GB对几十MB的差距倾斜实锤。第二步找到倾斜发生在哪个算子。Stage 5对应的算子就是ExchangeShuffle。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表假设目标表已经有分区字段dtspark.sql( INSERT OVERWRITE TABLE dws.user_wide PARTITION(dt20240101) 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语法你本来就会真正的深度在懂原理和会排查这两件事上。建议你拿到一个任务后别急着写完代码就跑先想清楚数据量级多大、关联的表多大、哪里可能倾斜、分区裁剪能不能生效。习惯养成了你的任务会跑得比别人稳很多。
返回列表