Spark 性能优化:代码写法先调,参数其次,最后加资源
Spark 任务慢,三个方向排查:
| 排查顺序 |
优化方向 |
典型手段 |
| 第一 |
代码层面 |
换算子、加缓存、广播变量 |
| 第二 |
参数层面 |
调并行度、调内存分配 |
| 第三 |
资源层面 |
加 Executor、加内存 |
不要一上来就加资源。先检查代码有没有写错,再看参数有没有配错,最后才考虑加机器
代码层面:90% 的性能问题出在这
1. 能用 reduceByKey,就别用 groupByKey
| 算子 |
行为 |
Shuffle 数据量 |
groupByKey |
所有数据 Shuffle 到 Reduce 端,再聚合 |
全量数据 |
reduceByKey |
Map 端先预聚合,再 Shuffle |
聚合后的数据 |
rdd.groupByKey().mapValues(_.sum)
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),每次都会重算。
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 都要用同一个大对象(比如字典表),用广播变量。
val dict = Map(...)
rdd.map(record => dict.get(record.key))
val bcDict = sc.broadcast(dict)
rdd.map(record => bcDict.value.get(record.key))