Spark 调优避坑指南:那些你以为优化了其实反向优化的操作
调优调到性能退化 30%,这事朱大喜这个月干了不止一次。如果你也中过招,这篇文章就是写给你的。
一、"直觉优化"是最危险的优化
Spark 的分布式计算模型和我们单机编程的直觉经常是反着来的。7 月份我参与了一个 ETL 项目优化,目标是 300GB Shuffle 数据,3 小时跑完 → 45 分钟。但在优化路径上,至少有 4 个"直觉上正确"的操作,实际效果是负的。
为什么 Spark 的分布式直觉和单机编程反着来?单机编程的世界里,内存越大越好、线程越多越好、提前过滤一定快——这些是线性思维。但在分布式系统里,每多一个分区就多一次网络传输和 Task 调度,每多一个 cache 就多一次序列化和 GC 扫描,每一步的代价都不是局部的。你给一个 Executor 加了 16G 内存,看起来单个 Task 跑得更从容了,但 G1GC 的 Full GC 扫描 16G 堆可能需要 30 秒——这 30 秒里整个 Executor 不干活,而 10 个 Executor 轮流 GC 的累积停滞时间可能比你"省下来的计算时间"还多。分布式性能优化的本质是理解每个操作的全局代价模型而非局部最优。
这篇文章不是给你看最佳实践的(那些网上太多了),而是把那些"做了更糟"的操作摊开给你看。
二、5 个反向优化的经典场景
反向优化 1:过度增加分区 — 2000 个分区比 200 个还慢
直觉:"数据量大,多分几个分区并行度更高,跑得快。"
真相:每个分区都是一个 Task,Task 调度本身有开销(序列化、网络通信、Task 启动)。当分区数远大于 executor 核数时,大量时间浪费在调度而非计算上。
// ❌ 反向优化:repartition(2000) 产生 2000 个小文件 df.repartition(2000) .write.mode("overwrite") .parquet("/data/output") // 实际效果:每个文件只有几 MB,元数据开销巨大 // 下游读取需要打开 2000 个文件,IOPS 打满 // ✅ 合理做法:根据数据量动态计算分区数 import org.apache.spark.sql.functions.spark_partition_id val totalSizeMB = 300000 // 300GB val targetPartitionMB = 128 // 目标每个分区 128MB(HDFS block 大小) val idealPartitions = (totalSizeMB / targetPartitionMB).toInt // ~2343 // 但分区数还要受 executor 核数约束 val totalCores = spark.conf.get("spark.executor.instances").toInt * spark.conf.get("spark.executor.cores").toInt // 比如 50 个 executor * 4 core = 200 val optimalPartitions = Math.min(idealPartitions, totalCores * 3) // 核数 × 2~3 df.coalesce(optimalPartitions) // coalesce 不会引起 shuffle .write.mode("overwrite") .parquet("/data/output")反向优化 2:滥用 Cache — 缓存了不该缓存的
直觉:"这个 DataFrame 后面还要用,cache 一下。"
真相:Cache 是拿内存换时间。如果 DataFrame 只被用一次,cache 纯粹是浪费内存 + 增加 GC 压力。而且 cache 本身的反序列化也有成本。
from pyspark.sql import SparkSession from pyspark.storagelevel import StorageLevel spark = SparkSession.builder.appName("CacheDecision").getOrCreate() # ❌ 反向优化:读一次就用的数据也 cache df_raw = spark.read.parquet("/data/events_202607") # 这个文件后续只做一次聚合就写出,cache 毫无意义 df_agg = df_raw.groupBy("user_id").agg({"amount": "sum"}) # 更糟的情况:cache 了会被多次 shuffle 的中间结果 df_intermediate = df_raw.join(df_dim, "user_id") # shuffle join df_intermediate.cache() # 缓存了一份 shuffle 后的数据 df_intermediate.count() # 触发 cache,数据在内存 # 但是!下一步还需要再做一次 shuffle groupBy df_result = df_intermediate.groupBy("category").sum() # 又触发 shuffle! # df_intermediate 马上就没用了,cache 白费 # ✅ 应该 cache 的场景: # 1. 被多次 Action 触发计算 # 2. 后续计算没有 shuffle(如 filter → 多次聚合) # 3. 迭代算法中重复使用的数据集 df_gold = df_raw.filter(df_raw['event_type'] == 'purchase') \ .select('user_id', 'amount', 'category') df_gold.cache() # 后续要用来做多种分析 # 用完后记得释放! df_gold.unpersist()反向优化 3:乱用 Python UDF — 性能杀手本人
直觉:"这个处理逻辑 Spark SQL 不好写,写个 Python UDF 吧。"
真相:PySpark UDF 需要把每行数据从 JVM 序列化到 Python 进程,处理完再序列化回去。这个开销在百万行级别就开始显现,上亿行直接不可接受。
为什么 Python UDF 的序列化开销是数量级的而不是百分比的?很多人以为序列化只是多花 10%-20% 的时间,但实际上是一个跨语言的 IPC 调用:每一行数据要从 JVM 的堆内存里拷贝出来,通过 Socket 或 Pipe 传给 Python 进程,Python 解析成 Python 对象,执行你的函数,再把结果序列化传回 JVM。1000 万行数据 = 1000 万次跨进程调用。相比之下,Spark SQL 的内置函数或 pandas UDF(Arrow 向量化批次传输)一次传一个数组给 Python,1000 万行可能只需要几百次调用。这不是 10% vs 100% 的差距,而是10-100 倍的差距。
from pyspark.sql.functions import udf, col from pyspark.sql.types import StringType # ❌ 反向优化:用 Python UDF 处理 1000 万行数据 @udf(returnType=StringType()) def categorize_amount_py(amount: float) -> str: """这个函数在 Python 进程执行,每行都要序列化/反序列化""" if amount < 100: return "low" elif amount < 1000: return "mid" else: return "high" # 1000 万行 × 序列化开销 = 灾难 df_bad = df.withColumn("amount_level", categorize_amount_py(col("amount"))) # ✅ 正确做法1:用 Spark SQL 内置函数 from pyspark.sql.functions import when df_good = df.withColumn( "amount_level", when(col("amount") < 100, "low") .when(col("amount") < 1000, "mid") .otherwise("high") ) # 全程在 JVM 内执行,0 序列化开销 # ✅ 正确做法2:如果逻辑复杂,用 pandas UDF(向量化) from pyspark.sql.functions import pandas_udf import pandas as pd @pandas_udf(returnType=StringType()) def categorize_amount_vectorized(amount_series: pd.Series) -> pd.Series: """pandas UDF 按批次处理,每批次是一个 pandas Series,大幅减少开销""" return pd.cut( amount_series, bins=[-float('inf'), 100, 1000, float('inf')], labels=['low', 'mid', 'high'] ) df_fast = df.withColumn("amount_level", categorize_amount_vectorized(col("amount"))) # pandas UDF 比普通 UDF 通常快 3-100 倍反向优化 4:一味增大 Executor 内存 — GC 反而更惨
直觉:"OOM 了?加内存!从 4G 加到 16G,爽!"
真相:Executor 内存越大,单次 GC 要扫描的堆越大。如果程序本身就有内存泄漏倾向(如大量 cache),加内存只是延缓 OOM,反而让 Full GC 的时间变得更长。
# ❌ 反向优化:盲目堆内存 --executor-memory 16G --conf spark.executor.memoryOverhead=2G # off-heap 配置太小 # Full GC 时:16G 堆 → 30s+ 的 stop-the-world # ✅ 合理配置:内存 + GC 策略配合 --executor-memory 8G --conf spark.executor.memoryOverhead=2G --conf spark.memory.fraction=0.6 # 60% 用于执行/存储(默认) --conf spark.memory.storageFraction=0.5 # 执行与存储各一半 --conf spark.executor.extraJavaOptions="-XX:+UseG1GC -XX:InitiatingHeapOccupancyPercent=35" # G1GC + 35% 触发并发标记,避免 Full GC反向优化 5:过早 Filter — 在错误的位置过滤
直觉:"早点过滤掉不用的数据,后面计算量就小了。"
真相:在 JOIN 之前做 filter,如果过滤条件在 JOIN 键对应的维度表上,实际上是在"预计算 JOIN",可能会丢失数据或需要额外的 broadcast。
// ❌ 反向优化:在 JOIN 前过度过滤 // 需求:找到 7 月下单的"VIP 用户"的行为数据 // 错误写法:先过滤 VIP 用户 val vipUsers = users.df.filter($"user_level" === "VIP") // 10 万 VIP val julyOrders = orders.df.filter($"order_date" === "2026-07-01") // 500 万订单 val result = julyOrders.join(vipUsers, "user_id") // 2 张大表的 shuffle join // ✅ 正确写法:先 JOIN 较小的维度条件,再做后续过滤 // 或者:用小表 broadcast import org.apache.spark.sql.functions.broadcast val resultOptimized = julyOrders .join(broadcast(vipUsers), "user_id") // VIP 用户表只有 10 万行,broadcast 不 shuffle .filter($"user_level" === "VIP") // 如果不过滤可以在 join 前做好三、一个真实的调优过程
这是我们 7 月实际优化的一个 ETL Job,从 3 小时 → 45 分钟的全过程:
关键参数变更:
# 调优前的 Spark 配置 — 典型"直觉配置" # spark.sql.shuffle.partitions = 2000 # 太多了 # spark.executor.memory = 16g # 代码里到处都是 .cache() # 调优后 spark.conf.set("spark.sql.shuffle.partitions", "400") # 核心数 × 2 spark.conf.set("spark.sql.adaptive.enabled", "true") # AQE 自适应 spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true") spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true") # 自动处理数据倾斜 spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "50MB") # 小表自动 broadcast四、调优决策清单
| 你想做的操作 | 先问自己 | 如果答案是 Yes,做;否则别做 |
|---|---|---|
| 增加分区 | 当前分区数 < 核数 × 2? | Yes → 增加 |
| cache | 这个数据会被触发 ≥ 3 次? | Yes → cache |
| Python UDF | 能用 Spark SQL 或 pandas UDF 替代? | Yes → 别用 Python UDF |
| 加内存 | 确认过 GC 日志不是瓶颈? | No → 先调 GC |
| 提前 filter | filter 之后 Join 的数据是不是显著变小了? | No → 先 Join 再 filter |
五、总结
🚨 踩坑提醒
repartition(N) 后数据倾斜可能更严重:
repartition是按 key 的 hash 分配数据到 N 个分区,如果你的 join key 本身就倾斜(比如某城市用户占 60%),不管 N 设多少,这些数据都会被哈希到同一个分区里,倾斜不会消失。此时应该用spark.sql.adaptive.skewJoin.enabled=true让 AQE 自动分裂倾斜分区,而不是调 N 的大小。coalesce 不能增加分区数:
coalesce(N)只能减少分区,不能增加。如果你误以为coalesce(200)可以把 50 个分区扩到 200 个,它其实什么都不做——因为 coalesce 没有 shuffle,无法重新分布数据。需要增加分区用repartition。AQE 的 coalescePartitions 和用户手动指定分区数会冲突:你设了
spark.sql.shuffle.partitions=400,但 AQE 的coalescePartitions会根据数据量自动合并小分区,最终实际分区数可能远小于 400。如果你的业务逻辑依赖"分区数不减少"(如分区内顺序有业务含义),需要关闭 AQE 的自动合并。
Spark 调优最反直觉的地方在于:"多"不一定好。多分区 = 多调度开销,多 cache = 多 GC 压力,多内存 = 多 GC 停顿,多 UDF = 多序列化成本。这个月最大的教训是:永远用 Spark UI 的指标说话,别用直觉说话。
看这几个指标就够了:
- Shuffle Read/Write Size:看数据倾斜
- Task Time 分布:看是否有长尾
- GC Time / Task Time:内存是否合理
- Input Size / Output Size:确认分区合理
数据不会骗人,直觉会。下个月继续优化 —— 目标是同一批 ETL 任务从 45 分钟压到 20 分钟。