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操作
- 当前工作流:
- 读取DDD数据为df_d、F1数据为df_1,执行UNION得到df_union
- 对phone_number字段执行explode,生成phone_number_exploded字段
- 先按phone_number_exploded分组,再按原phone_number分组,对hobbies、job_title字段应用
collect_set聚合 - 将最终结果写入S3
- 当前故障:
- 数据量增长后出现OOM,报错
ExecutorLostFailure和java.lang.OutOfMemoryError: Java heap space - 十亿级数据规模下出现数据倾斜,结果写入阶段最后一个任务长时间挂起
- 数据量增长后出现OOM,报错
解答
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、提前降低聚合数据量:
- 数据源预处理:
- 对F1-Fn等多格式文件,读取后执行
repartition或coalesce合并小分区,统一转换为Parquet格式存储到S3临时目录,降低后续UNION的开销。 - 对DDD数据提前过滤无效phone_number(如空数组、格式错误号码),减少explode后的数据量。
- 对F1-Fn等多格式文件,读取后执行
- 调整聚合逻辑:
该方案优势:仅触发两次Shuffle(分组聚合+Join),比原流程的两次分组Shuffle更高效;提前去重减少了聚合阶段的数据量,降低内存占用。// 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") ) - 存储优化:写入S3时按phone_number前缀做
partitionBy,同时设置spark.sql.parquet.compression.codec = "snappy"压缩数据,减少IO和存储开销。
3. AQE失效后的手动数据倾斜处理
当AQE无法自动解决倾斜时,可通过以下步骤干预:
- 识别倾斜键:通过Spark UI的Stage页面查看Shuffle Read数据分布,定位数据量远高于其他键的phone_number或phone_exploded值。
- 拆分倾斜数据单独处理:
- 分离倾斜与正常数据:
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:_*)) - 对倾斜数据加盐拆分分组:
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") ) - 合并倾斜与正常数据的聚合结果:
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
相关产品推荐
相关产品推荐

