
数据开发这行很多痛苦都来自同一句话这个任务怎么又跑了这么久我早先接手的 Spark 任务动不动就是十几分钟起步稍微加个复杂 join 就能奔着半小时去。后来我把在多个项目里沉淀下来的 Spark SQL/DataFrame 调优经验整理成一套方法代号就叫 Auron。这篇就把 Auron 的核心思路、诊断手段和具体优化步骤全部摊开讲清楚。它不是什么黑科技也不用换集群核心是先搞懂 Spark 慢在哪、为什么慢再动手改数据布局、改执行计划、改关键参数。适合正在被慢查询、资源排队、OOM 困扰的数据开发人员也适合刚接触 Spark 性能调优、想少走弯路的新手。1. 先弄明白Spark SQL 和 DataFrame 到底快在哪、慢在哪1.1 Spark SQL 和 DataFrame 是同一套引擎的两张脸很多人有个误区觉得 Spark SQL 写得复杂就慢换成 DataFrame API 写就快或者反过来。实际在 Spark 内部这两者会被翻译成同一种东西一棵经过 Catalyst 优化器处理的逻辑执行计划再转成物理执行计划最终走向 Tungsten 的二进制内存管理。换句话说DataFrame 和 Spark SQL 只是同一枚硬币的两面性能差异远远没有想象中那么大。那真正和慢相关的分水岭在哪里答案是你有没有给优化器留下足够的可优化空间。Spark SQL 的 Catalyst 优化器非常依赖规则比如谓词下推、列裁剪、常量折叠、分区裁剪。这些规则发挥作用的前提是查询里能清楚地表达“过滤条件是什么”“要读哪些列”“数据分布在哪些分区”。如果你在 SQL 里写了大量无法下推的 UDF或者在 DataFrame 链路里反复用 RDD 算子打断计划优化器就只能干瞪眼跑起来自然慢。我见过一个真实的例子同一份 2000 万行日志用 SQL 写与用 DataFrame API 写跑出来的执行计划几乎一模一样时间差在 3% 以内。反而是某位同事在 DataFrame 里为了取一个字段先 toJSON 再 parseJSON列裁剪彻底失效整张表几百个字段全扫了一遍任务直接慢了 8 倍。所以 Auron 的第一条原则就是先让执行计划变“好看”再谈参数调优。1.2 定位瓶颈不靠猜从 Spark UI 看这五个指标接到一个慢任务别急着改代码。先打开 Spark UI 的历史服务器找到对应 Application然后按顺序盯五个指标。第一个是 SQL 标签页里的 Physical Plan。这里能看到 Catalyst 优化后的执行计划重点看 Filter 是否被下推到 Scan 之前以及读取的列是否只有需要的那些。第二个是 Stage 列表里的 Duration 和 Shuffle Read Size。如果某个 Stage 的 Shuffle Read 特别大比如几个 GB那基本可以断定 shuffle 是主要瓶颈。第三个是每个 Task 的耗时分布。正常情况 Task 耗时应该比较均匀如果出现一两个 Task 耗时是平均值的几十倍这是典型的数据倾斜特征。第四个是 Executors 页面的 GC Time 和 Shuffle Write。GC 时间占比过高说明内存不够或者对象分配太频繁。第五个是 Event Timeline看任务到底卡在调度、计算还是 IO 等待。这五处看完慢在哪一层基本就有数了。Auron 的诊断流程里我坚持先看 UI 再动手原因很简单没有数据的调优就是盲人摸象你连瓶颈在哪个 Stage 都不知道调参只是在碰运气。2. 数据布局是第一层性能开关文件、分区和日期字段2.1 文件格式和分区设计为什么大家都在推 Parquet很多人优化 Spark 只盯着参数和 SQL却忽略了一个更基础的东西底层文件是怎么存的。数据布局决定扫描成本扫描成本往往占总耗时的大头。举个最简单的例子同一批数据存成 CSV 和存成 Parquet跑同样的聚合查询Parquet 经常能快 70% 以上因为列式存储天然支持列裁剪读聚合需要的列时其他列的数据根本不用碰。分区设计同样关键。一张日志表如果按 dt 字段做二级分区查询带上 dt 的等值或范围条件Spark 可以直接做分区裁剪跳过大量无关文件。反之如果分区字段乱来或者压根没分区每次查询都是全表扫描再怎么调参都救不回来。实际落地时我会建议按这个顺序检查数据布局第一文件格式是否列式存储第二表是否有合理分区字段第三单个数据文件大小是否在 128MB 上下第四是否有大量小于 16MB 的小文件。第四点非常容易被忽略但危害极大。小文件太多会导致 Spark 启动成千上万个 Task调度开销比计算本身还高。遇到这种情况先合并小文件再谈后面所有优化否则一切都白搭。2.2 日期字段处理加减、格式化和按月转存的隐藏成本日期字段是数据开发里最常用、也最容易写崩的一类字段。它看起来只是简单计算但在大数据量下一个函数的不当使用就可能让分区裁剪彻底失效。比如最常见的按月份过滤很多人会写where date_format(dt, yyyy-MM) 2024-06。这句话语法没错结果也对但 date_format 是一个表达式函数它会包裹住 dt 字段导致 Spark 无法把过滤条件下推到分区扫描层。原本只需要扫 6 月那一个分区的数据现在被迫扫全表所有分区再逐行做格式化性能差距可以到十倍以上。正确写法是范围过滤where dt 2024-06-01 and dt 2024-07-01这样分区裁剪才能真正生效。再比如日期加减。Spark 3 之后日期类型可以直接做算术运算比如order_date interval 1 year表示加一年order_date - interval 1 day表示减一天。但如果你还在用date_add(add_months(order_date, 12), -1)这种叠加写法不仅可读性差还会产生多余的表达式节点。对于大表上的日期转换尽量一次性算到位比如add_months(order_date, 12)就是标准的加一年不要用date_add(order_date, 365)后者在闰年场景下结果不对。日期按月转存也是高频场景。想把数据按月份落盘最常见做法是partitionBy(year, month)后用year(dt)和month(dt)生成分区列。这里有个坑如果你在写入之前对表做了全局排序比如orderBy(dt)那代价极高。正确做法是只对分区内做sortWithinPartitions既保证写入文件基本有序又避免一次全局 shuffle。关于日期格式转换还有一个细节from_unixtime/unix_timestamp这类函数会拿字符串去匹配一旦格式串写错Spark 不会报错而是返回 null后续计算全部出错。这类问题极难排查所以我建议在 ETL 入仓时就统一把日期存成标准date或timestamp类型不要在查询里反复做字符串解析。宁可入库时多花一点算力也不要让下游每个任务都做重复的类型转换。3. 合并、排序与聚合把 shuffle 压到最低3.1 数据合并先分清类型join、union 和广播的取舍数据合并是日常开发里最常见的操作也是性能问题的重灾区。合并之前先想清楚你要的是横向扩展列还是纵向追加行。横向扩展列那就是 join。join 的性能核心在于是否产生 shuffle。如果一个大表和一个小表做 join最理想的办法是广播小表让小表复制到每个 Executor 的内存里完全避免 shuffle。Spark 3 里只要小表小于spark.sql.autoBroadcastJoinThreshold默认会尝试广播阈值默认 10MB。但实际场景里这个阈值经常被改小或表刚好超一点导致走了 SortMergeJoin性能下降明显。Auron 的实践建议是对确定的小维度表显式加广播 hint比如/* BROADCAST(t2) */不要依赖自动判断。如果两个都是大表那 SortMergeJoin 基本躲不掉。此时唯一的优化空间是减少参与 shuffle 的数据量第一join 前先把过滤条件都压下去第二只 select 需要的列第三如果两个大表经常按同一个 key 关联可以考虑在建表时就做分桶让相同 key 落在同一个文件里触发 BucketJoin 避免 shuffle。分桶的代价是写入慢一点但对高频关联查询收益很大。纵向追加行也就是 union。这个操作本身不产生 shuffle两个 DataFrame 往下一拼接就行。但如果 union 之后还跟着一个全局orderBy那就会引发完整的 shuffle。很多报表场景其实不需要全局有序只要每个输出分区有序就足够所以优先用sortWithinPartitions。有些练习平台上的“数据合并”题通常要求把两个结构相同的表合并后去重或计算。实际项目里我会额外提醒一句用union时如果两张表列顺序不一致结果会错位。安全做法是用unionByName或者先显式 select 列名对齐。这个坑不影响性能但影响数据正确性而且错位问题非常隐蔽有时候跑完整个链路才发现某列数据全乱了返工成本极高。3.2 排序和聚合的正确姿势orderBy 不是你想用就能用排序是另一个容易“无脑写”的操作。一个全局orderBy底层一定触发全量 shuffle 加全量排序数据量一大耗时直线上升。写之前问自己三个问题这个排序是业务必需的吗是全局有序还是分区内有序就能满足能不能把排序推迟到最终结果集很小的时候再做如果结果集已经很小比如聚合之后的维度表那随便 orderBy 无所谓。但如果是对几百 GB 明细数据做全局排序就要慎重了。替代方案有两个一是sortWithinPartitions只对每个分区内部排序二是用repartition把数据按业务 key 分散到固定分区再在分区内排序。这两种方式都能保证“看起来有序”的写入效果同时避免全局大 shuffle。聚合计算同样有讲究。groupBy之后计算 sum/count/avgSpark 会先在每个 map 端做部分聚合再 shuffle 到 reduce 端做最终聚合这个机制叫部分聚合。它已经能压掉不少数据量。但如果你在聚合之前做了一次宽的 joinjoin 已经把数据膨胀了几倍那部分聚合的优势就被稀释了。所以 Auron 的建议是能先聚合再 join 的一定先聚合再 join能先过滤再聚合的一定先过滤再聚合。另一种常见的低效写法是用多个groupBy实现同一份数据的不同维度汇总比如既按用户算又按城市算又按品类算。三次 groupBy 意味着三次完整 shuffle。Spark 提供了rollup和cube可以用一次扫描加一次 shuffle 算出多组维度组合的聚合结果。数据量大时这个改写经常能带来 50% 以上的性能提升。数据倾斜是聚合和 join 场景里绕不开的问题。最常见的表现是 groupBy 之后某一个 key 的数据量特别大比如“不限城市”这种特殊值或者头部用户贡献了超多订单。此时单个 reduce task 会拖垮整个 Stage。Auron 用的标准解法是加盐给热点 key 加一个 0 到 N 的随机后缀把原来的一个大 key 拆成 N 个小 key第一轮聚合先处理分片后的 key第二轮把后缀去掉再做最终聚合。代价是多一轮聚合但能有效避免单点瓶颈。加盐的 N 不是越大越好一般按热点 key 的数据量除以目标分区大小来估算8 到 64 之间比较常见。4. 自适应执行与资源配置让 Spark 帮你调优4.1 开启 AQE动态合并分区、自动处理倾斜 join如果还在用 Spark 2.x那升级到 Spark 3.x 本身就是一次性能优化。最大红利来自 AQE全称 Adaptive Query Execution中文叫自适应查询执行。它的核心思想是把优化决策从编译期推迟到运行时。因为 Spark 真正跑起来之前谁也不知道每个 Stage 会输出多少数据、是否存在数据倾斜等到 shuffle 完成、数据量真实可见了再动态调整后续计划效果比任何静态调参都准。AQE 带来的第一个能力是动态合并 shuffle 分区。默认情况下spark.sql.shuffle.partitions是 200意思是所有 shuffle 之后的 Stage 都分 200 个分区。对于小数据量200 个分区会产生大量小 task浪费调度资源对于大数据量200 个分区又不够用。AQE 开启后Spark 会根据实际 shuffle 输出大小自动合并分区目标是把每个分区控制在合理的大小范围里。我在一个实际项目里见过订单表 join 维度表之后AQE 自动把 200 个分区合并成了 24 个reduce 阶段耗时缩短了 60%。第二个能力是动态切换 join 策略。如果一个表在运行时发现实际数据量远小于估算值AQE 会自动把它转成广播 join避免 SortMergeJoin 的 shuffle。第三个能力是动态优化倾斜 join。对于运行时发现的倾斜分区AQE 会自动拆分把一个倾斜分区拆成多个小任务并行处理并不需要你手动加盐。但注意AQE 的倾斜优化主要是针对 join 场景对 groupBy 这类聚合倾斜的处理能力有限所以前面讲的手动加盐仍然有用。开启 AQE 的配置非常简单在提交任务时加上这几行spark-submit \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ --conf spark.sql.adaptive.skewJoin.enabledtrue \ --conf spark.sql.adaptive.skewJoin.skewedPartitionFactor5 \ --conf spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes256MB \ application.py这几项配置里下面两个控制倾斜判断的敏感度。Factor 表示分区中位数乘以多少倍算倾斜Threshold 表示绝对值超过多少大小算倾斜。默认值偏保守如果你的场景确实存在已知倾斜可以适度调小 Threshold让 AQE 更早介入。4.2 资源参数与数据读取不是调越大越好资源参数是大家最爱调的东西但也是最容易“调坏”的地方。核心原则不是把每个 Executor 的内存拉满而是让资源配置和你的数据规模、计算复杂度匹配。第一个要算清楚的参数是 shuffle 分区数。一个经验法则是让每个 shuffle 分区的输出数据量在 100MB 到 200MB 之间。如果你的某个 Stage shuffle 输出总量是 20GB那分区数可以设置在 128 到 200 之间。直接套spark.sql.shuffle.partitions50这种拍脑袋值大概率会造成单个 Task 处理数据过多内存压力陡增。第二个要注意的是 Executor 的核数和内存比例。一个 Executor 分配 8 核那堆内内存最好不要低于 8GB否则每个 Task 能分到的内存太少GC 会非常频繁。相反如果每个 Executor 只给 2GB却配了 8 核那每个 Task 只有 256MB 内存排序和聚合很容易 OOM。常见的合理配置是单个 Executor 4 到 8 核内存按每核 2GB 到 3GB 来配。第三个是并行度和文件数的匹配。spark.sql.files.maxPartitionBytes默认 128MB如果单个文件特别大可以适当调小让一个文件被拆成多个分区并行读如果大量小文件则要反向调大避免每个小文件都启动一个 Task。关于数据读取还有一个常见痛点很多人习惯先用 pandas 读小文件做分析再转成 Spark DataFrame。在数据量小的时候这个流程没问题但一旦数据量上来了比如十几 GB还在用 pandas 读机器内存直接爆掉。Auron 的实践建议是数据量在单机内存可承载范围内用 pandas 没毛病一旦超过这个范围直接用spark.read去读分布式文件或者先用spark.createDataFrame(pandas_df)把 pandas 数据转成 Spark DataFrame 再处理。反过来从 Spark DataFrame 取少量结果到本地用.toPandas()没问题但如果全表调用内存必然会炸。记住一个原则Spark 是分布式的pandas 是单机的跨过这条线的数据量一律交给 Spark。5. 一次真实优化实录与问题速查表5.1 从 22 分钟到 3 分钟优化流程完整复盘有一个 ETL 任务让我印象很深一张 1500 万行的订单事实表 join 一张 300 万行的用户维度表然后按用户 ID 做 groupBy统计每人的订单总额和订单数最后按日期分区写出 Parquet。原始 SQL 大概长这样SELECT u.user_id, u.city, SUM(o.amount) AS total_amount, COUNT(o.order_id) AS order_cnt FROM orders o JOIN users u ON o.user_id u.user_id WHERE date_format(o.dt, yyyy-MM) 2024-06 GROUP BY u.user_id, u.city ORDER BY total_amount DESC这条 SQL 在集群上稳定跑 22 分钟。从 Spark UI 看过去最主要的问题有三个一是date_format(o.dt, yyyy-MM) 2024-06导致订单表分区裁剪全部失效整个 6 月数据前所有分区全被扫了一遍二是用户表 300 万行按默认阈值没有触发广播走了 SortMergeJoin大量数据在网络上传输三是最后的ORDER BY total_amount DESC是一次全局排序在全量聚合结果上其实没必要。优化分三步做。第一步把日期过滤改成范围过滤WHERE o.dt 2024-06-01 AND o.dt 2024-07-01第二步强制广播用户表SELECT /* BROADCAST(u) */ ...第三步去掉全局排序改成写入时分区内排序。写完后的 SQL 只跑了 8 分钟。再往下走打开 AQESpark 在 join 之后的聚合 Stage 自动把 200 个分区合并成了 36 个reduce 阶段耗时进一步下降。最终这个任务稳定在 3 分钟左右。全程没有加机器没有调大 Executor 内存改动的是数据访问方式和执行计划的选择。这正好说明了 Auron 的核心观点大部分慢查询不是资源不够而是没有用对方式让 Spark 干活。5.2 常见性能问题速查与解决路径问题现象常见原因优先排查手段解决方案路径任务整体很慢所有 Stage 都久数据量远超资源量或执行计划不合理查看 Physical Plan检查是否有全表扫描分区裁剪、列裁剪、过滤下推、优化数据布局某几个 Task 明显比兄弟 Task 慢数据倾斜看 Stage 页面 Task 耗时分布加盐拆 key或开启 AQE skewJoin热点 key 单独处理Shuffle Read 数据量巨大join 或 groupBy 前未充分过滤/聚合看每个 Stage 的 Shuffle Read Size先过滤、先聚合、再 join小表广播GC 时间占比高Executor 内存不足对象分配过多Executors 页面看 GC Time调整 Executor 内存与核数配比减少 shuffle 数据量输出大量小文件分区数设置不合理或大量 reduce task看输出文件个数与大小调整 shuffle 分区数写入前 coalesce写 Parquet 特别慢全局排序或序列化开销看 Stage 耗时改为 sortWithinPartitions检查序列化器配置这张表是 Auron 项目落地时整理出来的基本覆盖了我日常接触的 90% 问题。遇到性能瓶颈先对着表格找方向不要一上来就改spark.executor.memory很多时候问题根本不在这里。最后说一点个人体会。数据任务的性能和很多工程问题一样想当然的地方越多埋的雷就越多。Auron 这套方法的最大价值不是某一条具体的优化命令而是逼着我把“任务跑的慢”这个模糊感觉翻译成一个又一个可以观察、可以测量、可以对比的具体原因。先看执行计划再看 UI 指标再检查数据布局和 SQL 写法最后才动参数。按这个顺序走下来绝大多数 Spark SQL/DataFrame 性能问题都能找到明确答案。希望这篇整理能让你少踩一些我踩过的坑也欢迎在实施过程中随时交流你遇到的实际案例。