一、资源配置优化

核心目标:合理分配 CPU、内存、Executor 数量等资源,避免资源浪费或不足。

1. Executor 相关配置
  • 配置参数
    • --num-executors:Executor 数量(建议 50~100 个,过多会增加通信开销)。
    • --executor-cores:每个 Executor 的 CPU 核数(建议 2~5 核,核数过多会导致资源竞争)。
    • --executor-memory:每个 Executor 的内存(需预留部分给系统和缓存,如实际业务用 8G,配置 10G)。
  • 使用场景:适用于所有 Spark 作业,尤其是数据量大、计算密集型任务(如复杂聚合、Join)。
  • 示例(spark-submit):
    spark-submit \
      --num-executors 50 \
      --executor-cores 4 \
      --executor-memory 10g \
      --class com.example.MyJob \
      myjob.jar
    
2. Driver 配置
  • 配置参数
    • --driver-memory:Driver 内存(当有大量 collect () 操作或广播变量过大时需调大,建议 2~8G)。
  • 使用场景:作业中包含 collect()take() 等将数据拉取到 Driver 的操作,或广播变量较大时。
  • 示例
    spark-submit --driver-memory 4g ...
    
3. 动态资源调整
  • 配置参数
    • spark.dynamicAllocation.enabled=true:启用动态资源分配(根据作业负载自动增减 Executor)。
    • spark.shuffle.service.enabled=true:配合动态分配,避免 Shuffle 数据丢失。
  • 使用场景:资源紧张的集群,或作业负载波动大(如白天高峰、夜间低峰)。
  • 配置方式:在 spark-defaults.conf 中添加:
    spark.dynamicAllocation.enabled true
    spark.shuffle.service.enabled true
    

二、代码与算子优化

核心目标:避免低效算子,减少不必要的计算和数据传输。

1. 避免使用低效算子
  • 禁用算子collect()(将全量数据拉到 Driver,易 OOM)、foreach()(分布式执行但无返回值,不适合复杂逻辑)。
  • 替代方案
    • 用 take(n) 替代 collect()(只取前 n 条)。
    • 用 foreachPartition() 替代 foreach()(减少分区内对象创建开销,如数据库连接)。
  • 使用场景:需要获取部分结果或对每个分区做批量操作时。
  • 示例
    // 低效:foreach 每条数据创建连接
    rdd.foreach { x =>
      val conn = getConnection() // 每条数据创建连接,开销大
      conn.write(x)
    }
    
    // 高效:foreachPartition 每个分区创建一次连接
    rdd.foreachPartition { iter =>
      val conn = getConnection() // 分区内共享连接
      iter.foreach(conn.write)
    }
    
2. 选择合适的 Join 策略
  • Broadcast Join:小表广播到所有 Executor,避免 Shuffle(小表建议 < 100MB)。
    • 配置:spark.sql.autoBroadcastJoinThreshold=104857600(100MB,默认 10MB)。
    • 使用场景:大表与小表 Join(如事实表与维度表)。
  • Shuffle Join:适用于两个大表 Join,需确保资源充足(避免 OOM)。
  • 示例(强制广播):
    import org.apache.spark.sql.functions.broadcast
    val df = bigTable.join(broadcast(smallTable), "id") // 强制广播小表
    
3. 减少宽依赖(Shuffle 操作)
  • 宽依赖:如 groupByKey()reduceByKey()distinct() 等会触发 Shuffle,尽量用高效算子替代。
  • 优化方案
    • 用 reduceByKey() 替代 groupByKey()(本地预聚合,减少 Shuffle 数据量)。
    • 用 repartitionAndSortWithinPartitions() 替代 repartition() + sortByKey()(合并操作,减少一次 Shuffle)。
  • 示例
    // 低效:groupByKey 无预聚合,Shuffle 数据量大
    rdd.groupByKey().mapValues(_.sum)
    
    // 高效:reduceByKey 本地预聚合,减少 Shuffle 数据
    rdd.reduceByKey(_ + _)
    

三、数据处理优化

核心目标:减少数据量,提升读写效率。

1. 数据格式优化
  • 推荐格式:Parquet/ORC(列存、压缩,适合 Spark SQL 分析),避免 CSV/JSON(文本格式,解析慢、体积大)。
  • 使用场景:数据持久化(如中间结果存储)、表存储(Hive 表)。
  • 示例
    // 写入 Parquet 格式(自动压缩)
    df.write.mode("overwrite").parquet("hdfs://path/to/data")
    
2. 过滤与投影下推
  • 操作:尽早过滤数据(where()filter())和选择必要字段(select()),减少后续处理的数据量。
  • 使用场景:所有需要读取数据的作业,尤其是大表查询。
  • 示例
    // 低效:先全表读取再过滤
    df.select("*").filter("age > 18")
    
    // 高效:先过滤再选择字段(Spark 会优化为下推)
    df.filter("age > 18").select("id", "name")
    
3. 避免小文件问题
  • 问题:大量小文件会导致 Executor 任务过多,元数据管理开销大。
  • 解决方案
    • 写数据时合并小文件:df.repartition(10).write...(指定合理分区数)。
    • 读数据前合并:spark.read.option("mergeSchema", "true").parquet("path").coalesce(5)
  • 使用场景:数据写入 HDFS 或从 HDFS 读取时(如日志数据按天切割后产生大量小文件)。

四、Shuffle 优化

核心目标:减少 Shuffle 数据量,降低磁盘和网络开销。

1. Shuffle 并行度调整
  • 配置参数spark.sql.shuffle.partitions(Spark SQL 中 Shuffle 分区数,默认 200,可根据数据量调整)。
  • 原则:每个 Shuffle 分区数据量控制在 100MB~1GB 之间(避免过大导致 OOM,过小导致任务过多)。
  • 使用场景:Join、GroupBy 等触发 Shuffle 的操作,数据量远超默认分区处理能力时。
  • 示例
    spark.conf.set("spark.sql.shuffle.partitions", 500) // 调整为 500 个分区
    
2. Shuffle 数据压缩
  • 配置参数
    • spark.shuffle.compress=true:压缩 Shuffle 输出(默认 true,用 snappy 算法)。
    • spark.shuffle.spill.compress=true:压缩 Shuffle 溢写数据(默认 true)。
  • 使用场景:所有涉及 Shuffle 的作业,尤其数据量大时(减少磁盘 IO 和网络传输)。

五、内存管理优化

核心目标:合理分配内存用于缓存和计算,避免 OOM。

1. 内存比例调整
  • 配置参数
    • spark.memory.fraction:Executor 内存中用于执行和存储的比例(默认 0.6,剩余 0.4 给系统)。
    • spark.memory.storageFraction:存储内存占总内存的比例(默认 0.5,即执行和存储各占一半)。
  • 调优场景
    • 计算密集型(如复杂聚合):提高 spark.memory.fraction(如 0.7)。
    • 缓存密集型(如多次复用 RDD/DF):提高 spark.memory.storageFraction(如 0.6)。
2. 数据缓存策略
  • 操作:对复用的 RDD/DF 进行缓存(cache() 或 persist()),避免重复计算。
  • 缓存级别选择
    • MEMORY_ONLY:纯内存(最快,适合小数据)。
    • MEMORY_AND_DISK:内存不足时溢写到磁盘(适合中大数据)。
    • 避免 DISK_ONLY(速度慢)。
  • 示例
    val df = spark.read.parquet("path").persist(StorageLevel.MEMORY_AND_DISK) // 缓存
    df.createOrReplaceTempView("tmp")
    spark.sql("select count(*) from tmp") // 第一次计算并缓存
    spark.sql("select avg(age) from tmp") // 直接使用缓存,无需重算
    

六、其他优化

  1. 广播变量:将小数据集(如配置表)广播到所有 Executor,避免重复传输(SparkContext.broadcast())。

    • 使用场景:大表与小表 Join,或需要在每个 Executor 共享只读数据时。
  2. Checkpoint:对长依赖 RDD 做 Checkpoint(持久化到 HDFS),避免 lineage 过长导致失败后重算开销大。

    • 配置:spark.checkpoint.dir="hdfs://path/to/checkpoint",调用 rdd.checkpoint()
  3. JVM 调优:设置 Executor JVM 参数(如 --conf spark.executor.extraJavaOptions="-XX:+UseG1GC"),避免 GC 频繁。

总结

  • 资源配置是基础,需根据集群规模和作业类型调整。
  • 代码与算子优化是核心,减少 Shuffle 和低效操作。
  • 数据处理优化可从源头减少数据量,提升读写效率。
  • 调优需结合具体场景(如数据量、计算逻辑),通过 Spark UI(查看 Stage、Shuffle、GC 等指标)定位瓶颈后针对性优化。

更多推荐