# Spark 性能优化 10 个技巧
·
实战总结!让 Spark 任务速度提升 5-10 倍的核心优化方案
📌 前言
Spark 任务跑得慢?内存经常 OOM?数据倾斜让你怀疑人生?
这篇文章总结了我多年大数据开发中最实用的 10 个 Spark 性能优化技巧,每个技巧都配有代码示例和参数配置,直接抄作业就能用!
优化效果:
- ✅ 任务运行时间从 2 小时 → 15 分钟
- ✅ 内存使用降低 60%
- ✅ 数据倾斜问题彻底解决
🎯 技巧 1:合理设置分区数
问题场景
分区太少 → 任务并发度低,资源浪费
分区太多 → 任务调度开销大,小文件多
优化方案
// ❌ 错误:使用默认分区数
val df = spark.read.json("hdfs://path/to/data")
// ✅ 正确:根据数据量设置分区
val df = spark.read
.option("maxPartitionBytes", "134217728") // 128MB
.json("hdfs://path/to/data")
// 或者手动 repartition
val df = df.repartition(200) // 根据集群核心数调整
参数建议
| 数据量 | 建议分区数 | maxPartitionBytes |
|---|---|---|
| < 10GB | 50-100 | 256MB |
| 10-100GB | 200-500 | 128MB |
| > 100GB | 1000+ | 64MB |
核心公式
分区数 = 总数据量 / 128MB
理想每分区数据量:100-200MB
🎯 技巧 2:使用广播变量(Broadcast)
问题场景
小表 Join 大表时,Shuffle 开销巨大
优化方案
// ❌ 错误:普通 Join 触发 Shuffle
val result = largeDF.join(smallDF, "id")
// ✅ 正确:广播小表,避免 Shuffle
import org.apache.spark.sql.functions.broadcast
val result = largeDF.join(broadcast(smallDF), "id")
适用条件
- 小表大小 < 10MB(默认阈值)
- 小表可以完整加载到内存
参数调整
// 调整广播阈值(默认 10MB)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "52428800") // 50MB
性能对比
| 场景 | 普通 Join | 广播 Join |
|---|---|---|
| 1 亿 × 1 万 | 30 分钟 | 2 分钟 |
| Shuffle 数据量 | 50GB | 0GB |
🎯 技巧 3:选择合适的缓存策略
问题场景
重复使用 DataFrame 时,每次都重新计算
优化方案
// ❌ 错误:不缓存,重复计算
val df = spark.read.parquet("hdfs://path")
df.count() // 计算 1
df.filter(...).show() // 重新计算
// ✅ 正确:使用 cache()
val df = spark.read.parquet("hdfs://path").cache()
df.count() // 计算并缓存
df.filter(...).show() // 直接使用缓存
// ✅ 进阶:使用 persist() 指定存储级别
import org.apache.spark.storage.StorageLevel
df.persist(StorageLevel.MEMORY_AND_DISK)
存储级别选择
| 场景 | 推荐级别 | 说明 |
|---|---|---|
| 内存充足 | MEMORY_ONLY |
最快 |
| 内存紧张 | MEMORY_AND_DISK |
溢出到磁盘 |
| 容错要求高 | MEMORY_AND_DISK_2 |
2 副本 |
| 序列化存储 | MEMORY_ONLY_SER |
节省空间 |
缓存管理
// 查看缓存
spark.catalog.listTables()
// 手动释放缓存
df.unpersist()
// 清理所有缓存
spark.catalog.clearCache()
🎯 技巧 4:解决数据倾斜(Salt 加盐法)
问题场景
某个 Key 数据量过大,导致单个 Task 运行时间远超其他
优化方案
// ❌ 错误:直接聚合,倾斜严重
val result = df.groupBy("user_id").agg(...)
// ✅ 正确:加盐分散,二次聚合
import org.apache.spark.sql.functions._
// 1. 添加随机盐值(0-9)
val saltedDF = df.withColumn("salt", (rand() * 10).cast("int"))
// 2. 第一次聚合(按 user_id + salt)
val partialDF = saltedDF
.groupBy("user_id", "salt")
.agg(sum("amount").as("partial_amount"))
// 3. 第二次聚合(按 user_id)
val result = partialDF
.groupBy("user_id")
.agg(sum("partial_amount").as("total_amount"))
适用场景
- Key 分布极度不均(如热点用户)
- 聚合后数据量大幅减少
盐值数量建议
盐值数 = 倾斜 Key 数据量 / 平均数据量
通常设置 10-100
🎯 技巧 5:解决数据倾斜(自定义分区器)
问题场景
Join 操作时,某个 Key 的数据集中在一个分区
优化方案
// ✅ 方案:自定义分区器,打散倾斜 Key
import org.apache.spark.Partitioner
class SkewPartitioner(numPartitions: Int, skewKeys: Set[String])
extends Partitioner {
override def numPartitions: Int = numPartitions
override def getPartition(key: Any): Int = {
val keyStr = key.toString
if (skewKeys.contains(keyStr)) {
// 倾斜 Key 分散到多个分区
scala.util.Random.nextInt(numPartitions / 2)
} else {
// 正常 Key 均匀分布
keyStr.hashCode % numPartitions
}
}
}
// 使用示例
val skewKeys = Set("hot_user_001", "hot_user_002")
val partitionedRDD = rdd.partitionBy(
new SkewPartitioner(200, skewKeys)
)
🎯 技巧 6:优化序列化方式
问题场景
默认 Java 序列化慢且占用空间大
优化方案
// ✅ 使用 Kryo 序列化
val conf = new SparkConf()
.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.set("spark.kryo.registrationRequired", "false") // 开发环境
.set("spark.kryo.unsafe", "true") // 性能优先
// 注册常用类(生产环境建议)
val conf = new SparkConf()
.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.set("spark.kryo.registrationRequired", "true")
.set("spark.kryo.classesToRegister",
"com.example.MyClass,com.example.MyUDF")
性能对比
| 序列化方式 | 速度 | 空间占用 | 适用场景 |
|---|---|---|---|
| Java 默认 | 慢 | 大 | 兼容性要求高 |
| Kryo | 快 3-5 倍 | 小 50% | 生产环境 |
🎯 技巧 7:调整 Shuffle 参数
问题场景
Shuffle 阶段性能瓶颈,磁盘 IO 过高
优化方案
// 核心参数配置
spark.conf.set("spark.sql.shuffle.partitions", "200") // 根据数据量调整
spark.conf.set("spark.default.parallelism", "200")
// 启用自适应查询执行(AQE)- Spark 3.0+
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
// 调整 Shuffle 缓冲区
spark.conf.set("spark.shuffle.file.buffer", "1024k") // 默认 32k
spark.conf.set("spark.reducer.maxSizeInFlight", "128m") // 默认 48m
AQE 自动优化
✅ 自动合并小分区
✅ 自动优化 Join 策略
✅ 自动处理数据倾斜
参数建议
| 集群规模 | shuffle.partitions | parallelism |
|---|---|---|
| 10 核心 | 50-100 | 50 |
| 50 核心 | 200-300 | 200 |
| 100+ 核心 | 500-1000 | 500 |
🎯 技巧 8:避免使用 UDF,使用内置函数
问题场景
UDF 无法优化,且序列化开销大
优化方案
// ❌ 错误:使用 UDF
import org.apache.spark.sql.functions.udf
val toUpperUDF = udf((s: String) => s.toUpperCase)
val result = df.withColumn("upper_name", toUpperUDF($"name"))
// ✅ 正确:使用内置函数
import org.apache.spark.sql.functions.upper
val result = df.withColumn("upper_name", upper($"name"))
性能对比
| 方式 | 执行计划 | 性能 |
|---|---|---|
| UDF | 黑盒,无法优化 | 慢 |
| 内置函数 | Catalyst 优化 | 快 2-3 倍 |
必须用 UDF 时的优化
// 使用 Pandas UDF(向量化)
import org.apache.spark.sql.expressions.pandas_udf
@pandas_udf(StringType)
def to_upper(v: pd.Series) -> pd.Series:
return v.str.upper()
🎯 技巧 9:优化文件读取
问题场景
读取大量小文件,NameNode 压力大,启动慢
优化方案
// ❌ 错误:直接读取大量小文件
val df = spark.read.parquet("hdfs://path/to/many/small/files")
// ✅ 正确:合并小文件后读取
// 方案 1:使用 wholeTextFiles + 合并
val files = sc.wholeTextFiles("hdfs://path")
val combined = files.values.reduce(_ + _)
// 方案 2:预处理合并(推荐)
// 使用 hadoop 命令合并
// hdfs dfs -getmerge /source /local/path
// hdfs dfs -put /local/path/combined.parquet /dest
// 方案 3:调整分区发现
val df = spark.read
.option("recursiveFileLookup", "true")
.parquet("hdfs://path")
文件格式选择
| 格式 | 压缩比 | 读取速度 | 适用场景 |
|---|---|---|---|
| Parquet | 高 | 快 | 列式查询 |
| ORC | 高 | 快 | Hive 兼容 |
| JSON | 低 | 慢 | 半结构化 |
| CSV | 低 | 中 | 数据交换 |
读取优化参数
spark.conf.set("spark.sql.files.maxPartitionBytes", "134217728")
spark.conf.set("spark.sql.files.openCostInBytes", "4194304")
🎯 技巧 10:监控与调优
问题场景
不知道瓶颈在哪里,盲目调优
优化方案
// 1. 启用 Spark UI
// spark-submit --conf spark.ui.port=4040
// 2. 查看执行计划
df.explain(true) // 查看详细物理计划
// 3. 监控指标
spark.sparkContext.listenerBus.post(SparkListenerTaskEnd(...))
// 4. 使用 Spark History Server
// 查看历史任务执行情况
关键指标监控
| 指标 | 正常范围 | 异常处理 |
|---|---|---|
| Task 运行时间 | < 1 分钟 | 检查数据倾斜 |
| GC 时间占比 | < 10% | 增加内存 |
| Shuffle 读写比 | < 5 | 优化 Join |
| 反压 | 无 | 调整批次大小 |
调优流程
1. 查看 Spark UI → 找到慢 Task
2. 分析执行计划 → 确定瓶颈
3. 针对性优化 → 参数/代码
4. 对比验证 → 确认效果
📊 综合优化案例
场景:电商订单分析
// 优化前:运行 2 小时
val orders = spark.read.parquet("hdfs://orders")
val users = spark.read.parquet("hdfs://users")
val result = orders
.join(users, "user_id")
.groupBy("user_id", "category")
.agg(sum("amount").as("total_amount"))
.orderBy($"total_amount".desc)
// 优化后:运行 15 分钟
val conf = new SparkConf()
.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.set("spark.sql.shuffle.partitions", "200")
.set("spark.sql.adaptive.enabled", "true")
val orders = spark.read
.option("maxPartitionBytes", "134217728")
.parquet("hdfs://orders")
.cache()
val users = spark.read.parquet("hdfs://users")
val result = orders
.join(broadcast(users), "user_id") // 广播小表
.groupBy("user_id", "category")
.agg(sum("amount").as("total_amount"))
.orderBy($"total_amount".desc)
result.cache() // 缓存结果
优化效果对比
| 指标 | 优化前 | 优化后 | 提升 |
|---|---|---|---|
| 运行时间 | 120 分钟 | 15 分钟 | 8 倍 |
| Shuffle 数据 | 500GB | 50GB | 10 倍 |
| 内存使用 | 80GB | 30GB | 63% ↓ |
📋 参数速查表
# 序列化
spark.serializer=org.apache.spark.serializer.KryoSerializer
# 并行度
spark.default.parallelism=200
spark.sql.shuffle.partitions=200
# 内存
spark.executor.memory=4g
spark.driver.memory=2g
spark.memory.fraction=0.6
# Shuffle
spark.shuffle.file.buffer=1024k
spark.reducer.maxSizeInFlight=128m
# AQE(Spark 3.0+)
spark.sql.adaptive.enabled=true
spark.sql.adaptive.coalescePartitions.enabled=true
spark.sql.adaptive.skewJoin.enabled=true
# 广播
spark.sql.autoBroadcastJoinThreshold=52428800
# 文件
spark.sql.files.maxPartitionBytes=134217728
✅ 总结
| 技巧 | 适用场景 | 提升幅度 |
|---|---|---|
| 合理分区 | 所有任务 | 2-3 倍 |
| 广播变量 | 小表 Join | 5-10 倍 |
| 缓存策略 | 重复计算 | 2-5 倍 |
| 数据倾斜 | 热点 Key | 10 倍 + |
| Kryo 序列化 | 大数据量 | 2-3 倍 |
| Shuffle 调优 | 聚合/Join | 2-4 倍 |
| 避免 UDF | 复杂转换 | 2-3 倍 |
| 文件优化 | 小文件多 | 3-5 倍 |
核心原则:
- 先监控,再优化
- 先代码,后参数
- 先局部,后全局
🔗 参考资料
💬 你在 Spark 优化中遇到过哪些坑?欢迎评论区交流!
📌 下一篇预告:《数据仓库分层设计指南》—— 从 0 搭建企业级数仓架构
更多推荐



所有评论(0)