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(_ + _)