You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark十亿级数据Enrichment作业OOM及数据倾斜问题优化咨询

数据Enrichment管道内存溢出与数据倾斜问题解决方案

问题背景

  • 管道目标:处理至少10TB数据,当前仅处理180GB(不足目标2%)即触发内存溢出(OOM)
  • 集群配置:10节点集群,单节点配置为16CPU、64GB内存、256GB磁盘;总计150核、570GB内存;处理结果存储于S3(Parquet格式)
  • Spark配置:deploy.mode = "client",master = "spark://192.168.x.x:7077",default.parallelism = 450,sql.shuffle.partitions = 450,已启用自适应查询执行(AQE)
  • 业务场景:
    • 数据源:S3上180GB Parquet格式的DDD数据,以及多格式(CSV/TXT/XLSX)的F1-Fn文件(大小覆盖MB级到50GB)
    • 数据Schema:包含name、phone_number(Array类型)等30+列,需按phone_number分组完成Enrichment操作
  • 当前工作流:
    1. 读取DDD数据为df_d、F1数据为df_1,执行UNION得到df_union
    2. 对phone_number字段执行explode,生成phone_number_exploded字段
    3. 先按phone_number_exploded分组,再按原phone_number分组,对hobbies、job_title字段应用collect_set聚合
    4. 将最终结果写入S3
  • 当前故障:
    • 数据量增长后出现OOM,报错ExecutorLostFailure和java.lang.OutOfMemoryError: Java heap space
    • 十亿级数据规模下出现数据倾斜,结果写入阶段最后一个任务长时间挂起

解答

1. 当前流程的核心问题

  • 重复分组导致冗余Shuffle:两次连续分组操作会触发两次Shuffle,尤其是explode后数据量倍数级增长,第二次分组会进一步放大内存和IO开销,直接加剧OOM风险。
  • collect_set聚合的内存压力:若某个phone_number对应的分组数据量极大,Executor需要在内存中存储该分组下大量去重后的hobbies、job_title数据,极易触发堆内存溢出。
  • UNION操作的隐式风险:若DDD与F类数据的Schema未严格对齐(如列顺序、类型不一致),UNION可能导致数据异常;同时F类多为小文件数据源,未合并分区会导致后续任务并行度不均,加剧数据倾斜。
  • Client部署模式的瓶颈:deploy.mode = "client"下Driver运行在提交节点,而非集群内部,处理大数量级数据时,Driver需接收大量任务状态信息,易成为性能瓶颈甚至触发Driver端OOM。

2. 优化后的工作流(用Join替代重复分组)

核心思路是减少冗余Shuffle、提前降低聚合数据量:

  1. 数据源预处理:
    • 对F1-Fn等多格式文件,读取后执行repartition或coalesce合并小分区,统一转换为Parquet格式存储到S3临时目录,降低后续UNION的开销。
    • 对DDD数据提前过滤无效phone_number(如空数组、格式错误号码),减少explode后的数据量。
  2. 调整聚合逻辑:
    // 1. 用unionByName合并所有数据源,确保Schema对齐
    val df_all = df_d.unionByName(df_1).unionByName(df_2)...
    
    // 2. explode phone_number并提前去重,减少聚合压力
    val df_exploded = df_all.select($"name", $"phone_number", $"hobbies", $"job_title")
      .withColumn("phone_exploded", explode($"phone_number"))
      .dropDuplicates("phone_exploded", "hobbies", "job_title")
    
    // 3. 按拆分后的号码分组聚合
    val df_agg = df_exploded.groupBy("phone_exploded")
      .agg(collect_set("hobbies").alias("hobbies_set"), collect_set("job_title").alias("job_title_set"))
    
    // 4. 关联回原数据,按原phone_number合并结果
    val df_result = df_all.join(df_agg, array_contains(df_all("phone_number"), df_agg("phone_exploded")), "left")
      .groupBy("name", "phone_number")
      .agg(
        collect_set("hobbies_set").alias("temp_hobbies"),
        collect_set("job_title_set").alias("temp_job_title")
      )
    
    // 5. 扁平化并最终去重(按需)
    val df_final = df_result.select(
      $"name", $"phone_number",
      array_distinct(flatten($"temp_hobbies")).alias("hobbies"),
      array_distinct(flatten($"temp_job_title")).alias("job_title")
    )
    
    该方案优势:仅触发两次Shuffle(分组聚合+Join),比原流程的两次分组Shuffle更高效;提前去重减少了聚合阶段的数据量,降低内存占用。
  3. 存储优化:写入S3时按phone_number前缀做partitionBy,同时设置spark.sql.parquet.compression.codec = "snappy"压缩数据,减少IO和存储开销。

3. AQE失效后的手动数据倾斜处理

当AQE无法自动解决倾斜时,可通过以下步骤干预:

  • 识别倾斜键:通过Spark UI的Stage页面查看Shuffle Read数据分布,定位数据量远高于其他键的phone_number或phone_exploded值。
  • 拆分倾斜数据单独处理:
    1. 分离倾斜与正常数据:
      val skewed_keys = Seq("138xxxxxxx", "139xxxxxxx") // 从UI识别的倾斜键
      val df_skewed = df_exploded.where($"phone_exploded".isin(skewed_keys:_*))
      val df_normal = df_exploded.where(!$"phone_exploded".isin(skewed_keys:_*))
      
    2. 对倾斜数据加盐拆分分组:
      val df_skewed_salted = df_skewed.withColumn("salt", (rand() * 10).cast(IntegerType))
        .withColumn("phone_salted", concat($"phone_exploded", lit("_"), $"salt"))
      val df_skewed_agg = df_skewed_salted.groupBy("phone_salted")
        .agg(collect_set("hobbies").alias("hobbies_set"), collect_set("job_title").alias("job_title_set"))
        .withColumn("phone_exploded", split($"phone_salted", "_")(0))
        .groupBy("phone_exploded")
        .agg(
          array_distinct(flatten(collect_set("hobbies_set"))).alias("hobbies_set"),
          array_distinct(flatten(collect_set("job_title_set"))).alias("job_title_set")
        )
      
    3. 合并倾斜与正常数据的聚合结果:
      val df_agg_total = df_normal.groupBy("phone_exploded")
        .agg(collect_set("hobbies").alias("hobbies_set"), collect_set("job_title").alias("job_title_set"))
        .unionByName(df_skewed_agg)
      
  • 调整Spark资源配置:
    • 增大executor.memory(如单Executor分配32GB,结合单节点64GB内存,配置2个Executor/节点),同时设置executor.memoryOverhead为内存的20%-30%,避免堆外内存溢出。
    • 调大sql.shuffle.partitions(如从450调整为1000-2000),确保每个Shuffle分区数据量控制在100MB以内。
  • 切换部署模式:将deploy.mode改为cluster,让Driver运行在集群节点上,利用集群内存资源,避免Driver成为瓶颈。

内容的提问来源于stack exchange,提问作者yaoviametepe

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.13 09:35:21