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 个还慢直觉数据量大多分几个分区并行度更高跑得快。真相每个分区都是一个 TaskTask 调度本身有开销序列化、网络通信、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 // 目标每个分区 128MBHDFS 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 UDFArrow 向量化批次传输一次传一个数组给 Python1000 万行可能只需要几百次调用。这不是 10% vs 100% 的差距而是10-100 倍的差距。from pyspark.sql.functions import udf, col from pyspark.sql.types import StringType # ❌ 反向优化用 Python UDF 处理 1000 万行数据 udf(returnTypeStringType()) 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(returnTypeStringType()) 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.memoryOverhead2G # off-heap 配置太小 # Full GC 时16G 堆 → 30s 的 stop-the-world # ✅ 合理配置内存 GC 策略配合 --executor-memory 8G --conf spark.executor.memoryOverhead2G --conf spark.memory.fraction0.6 # 60% 用于执行/存储默认 --conf spark.memory.storageFraction0.5 # 执行与存储各一半 --conf spark.executor.extraJavaOptions-XX:UseG1GC -XX:InitiatingHeapOccupancyPercent35 # 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做否则别做增加分区当前分区数 核数 × 2Yes → 增加cache这个数据会被触发 ≥ 3 次Yes → cachePython UDF能用 Spark SQL 或 pandas UDF 替代Yes → 别用 Python UDF加内存确认过 GC 日志不是瓶颈No → 先调 GC提前 filterfilter 之后 Join 的数据是不是显著变小了No → 先 Join 再 filter五、总结 踩坑提醒repartition(N) 后数据倾斜可能更严重repartition是按 key 的 hash 分配数据到 N 个分区如果你的 join key 本身就倾斜比如某城市用户占 60%不管 N 设多少这些数据都会被哈希到同一个分区里倾斜不会消失。此时应该用spark.sql.adaptive.skewJoin.enabledtrue让 AQE 自动分裂倾斜分区而不是调 N 的大小。coalesce 不能增加分区数coalesce(N)只能减少分区不能增加。如果你误以为coalesce(200)可以把 50 个分区扩到 200 个它其实什么都不做——因为 coalesce 没有 shuffle无法重新分布数据。需要增加分区用repartition。AQE 的 coalescePartitions 和用户手动指定分区数会冲突你设了spark.sql.shuffle.partitions400但 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 分钟。