spark优化
Spark 性能优化:代码写法先调,参数其次,最后加资源
Spark 任务慢,三个方向排查:
| 排查顺序 | 优化方向 | 典型手段 |
|---|---|---|
| 第一 | 代码层面 | 换算子、加缓存、广播变量 |
| 第二 | 参数层面 | 调并行度、调内存分配 |
| 第三 | 资源层面 | 加 Executor、加内存 |
不要一上来就加资源。先检查代码有没有写错,再看参数有没有配错,最后才考虑加机器
代码层面:90% 的性能问题出在这
1. 能用 reduceByKey,就别用 groupByKey
| 算子 | 行为 | Shuffle 数据量 |
|---|---|---|
groupByKey |
所有数据 Shuffle 到 Reduce 端,再聚合 | 全量数据 |
reduceByKey |
Map 端先预聚合,再 Shuffle | 聚合后的数据 |
// 低效:所有数据都 Shuffle
rdd.groupByKey().mapValues(_.sum)
// 高效:Map 端预聚合
rdd.reduceByKey(_ + _)
2. 能用 mapPartitions,就别用 map
map 每条数据调用一次函数,mapPartitions 每个分区调用一次。
// 错误:每条数据创建一次数据库连接
rdd.map(record => {
val conn = getConnection() // 重复创建
process(record, conn)
})
// 正确:每个分区创建一次连接
rdd.mapPartitions(iter => {
val conn = getConnection() // 一个分区一次
iter.map(record => process(record, conn))
})
3. 重复使用的 RDD 要缓存
同一个 RDD 被多次使用(比如多个 Action),每次都会重算。
// 错误: 两次 action 都重算
val rdd = sc.textFile("hdfs://data.txt").flatMap(...).map(...)
rdd.count()
rdd.collect()
// 正确: 中间结果缓存
val rdd = sc.textFile("hdfs://data.txt").flatMap(...).map(...).cache()
rdd.count()
rdd.collect()
4. 广播大变量,不要每个 Task 传一份
如果每个 Task 都要用同一个大对象(比如字典表),用广播变量。
// 错误: 每个 Task 都发一份(100 个 Task = 100 份)
val dict = Map(...) // 大对象
rdd.map(record => dict.get(record.key))
// 正确: 广播到每个 Executor(10 个 Executor = 10 份)
val bcDict = sc.broadcast(dict)
rdd.map(record => bcDict.value.get(record.key))
参数层面:调并行度、内存、GC
代码优化完了还慢,再看参数。
1. 并行度:每个分区 128MB-256MB
| 参数 | 默认值 | 说明 |
|---|---|---|
spark.default.parallelism |
集群总核数 | Shuffle 操作的分区数 |
spark.sql.shuffle.partitions |
200 | Spark SQL 的 Shuffle 分区数 |
判断标准: 每个 Task 处理 128MB-256MB 数据最合适。太小调度开销大,太大 OOM 风险高。
// 手动调分区数
rdd.reduceByKey(_ + _, 200) // 指定 200 个分区
// 重分区
rdd.repartition(200) // 增加分区(会 Shuffle)
rdd.coalesce(50) // 减少分区(不 Shuffle)
2. 内存分配
| 参数 | 默认值 | 建议 |
|---|---|---|
spark.executor.memory |
1GB | 根据数据量调到 8-32GB |
spark.driver.memory |
1GB | collect() 结果大就调高 |
spark.memory.fraction |
0.6 | 内存密集型调到 0.7-0.8 |
3. GC 优化
大内存 Executor 容易出现 GC 停顿,换 G1 收集器:
spark-env.sh 里加
export SPARK_JAVA_OPTS="-XX:+UseG1GC -XX:MaxGCPauseMillis=200"
资源层面:加 Executor、加内存、加核
代码和参数都调完了还不够,才考虑加资源。
提交任务时的资源配置:
spark-submit \
--num-executors 10 \
--executor-cores 4 \
--executor-memory 16g \
--driver-memory 4g \
--conf spark.default.parallelism=80 \
myapp.jar
资源配置原则:
| 参数 | 原则 |
|---|---|
| Executor 数量 | 集群总核数 ÷ 每 Executor 核数 |
| 每 Executor 核数 | 2-5 核,太多 GC 压力大 |
| 每 Executor 内存 | 核数 × 4-8GB |
| 并行度 | 总核数 × 2-3 倍 |
数据倾斜:卡在 99% 的常见原因
现象: 大部分 Task 完成了,一两个 Task 卡在 99% 不动。
诊断: Spark UI → Stage → 看每个 Task 的 Shuffle Read 数据量,某个 Task 读的数据量是别人的几十倍,就是倾斜。
解决方案:
1. 加盐分散(最常用)
// 对倾斜 Key 加随机前缀,分散到多个 Task 聚合
val skewedKey = "hot_key"
val result = rdd
.map { case (k, v) =>
if (k == skewedKey) {
(s"${k}_${Random.nextInt(10)}", v) // 加盐,分散到 10 个 Task
} else {
(k, v)
}
}
.reduceByKey(_ + _) // 第一轮聚合(分散)
.map { case (k, v) =>
(k.split("_")(0), v) // 去盐
}
.reduceByKey(_ + _) // 第二轮聚合(汇总)
2. 倾斜 Key 单独处理
val skewedKey = "hot_key"
val skewedRDD = rdd.filter(_._1 == skewedKey)
val normalRDD = rdd.filter(_._1 != skewedKey)
// 正常聚合
val normalResult = normalRDD.reduceByKey(_ + _)
// 倾斜部分加盐聚合
val skewedResult = skewedRDD
.map { case (k, v) => (s"${k}_${Random.nextInt(10)}", v) }
.reduceByKey(_ + _)
.map { case (k, v) => (k.split("_")(0), v) }
.reduceByKey(_ + _)
// 合并
val finalResult = normalResult.union(skewedResult).reduceByKey(_ + _)