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

Spark写入S3 109GB数据耗时7小时,寻求配置优化方案

Spark作业写入S3耗时过长的优化分析与配置方案

问题背景

我在EMR集群上运行Spark作业,需求是读取14个S3对象并进行左连接,已尝试广播连接、持久化中间结果、调整Spark配置等手段,但写入S3仍耗时7小时,寻求优化方案。

当前实现细节

输入数据处理(持久化输入数据)

读取Parquet文件,选择所需列并去重后持久化:

df1 = spark.read.parquet("\file1").select(required columns).distinct() --173063条记录
df1.persist()
df2 = spark.read.parquet("\file2").select(required columns).distinct() --140424条记录
df2.persist()
df3 = spark.read.parquet("\file3").select(required columns).distinct() --4255条记录
df3.persist()
df4 = spark.read.parquet("\file4").select(required columns).distinct() --475750条记录
df4.persist()
df5 = spark.read.parquet("\file5").select(required columns).distinct() --95778条记录
df5.persist()
df6 = spark.read.parquet("\file6").select(required columns).distinct() --63464条记录
df6.persist()
df7 = spark.read.parquet("\file7").select(required columns).distinct() --2250条记录
df7.persist()
df8 = spark.read.parquet("\file8").select(required columns).distinct() --360条记录
df8.persist()
df9 = spark.read.parquet("\file9").select(required columns).distinct() --202297条记录
df9.persist()
df10 = spark.read.parquet("\file110").select(required columns).distinct() --134573条记录
df10.persist()
df11 = spark.read.parquet("\file11").select(required columns).distinct() --125541条记录
df11.persist()
df12 = spark.read.parquet("\file12").select(required columns).distinct() --202条记录
df12.persist()
df13 = spark.read.parquet("\file13").select(required columns).distinct() --207条记录
df13.persist()
df14 = spark.read.parquet("\file14").select(required columns).distinct() --21124条记录
df14.persist()

连接逻辑实现

依次执行广播左连接,每次连接后删除无用列并持久化中间结果:

join1_df = df1.join(broadcast(df2),left).drop(not required columns used in join)
join1_df.persist()
join2_df = join1_df.join(broadcast(df3),left).drop(not required columns used in join)
join2_df.persist()
join3_df = join2_df.join(broadcast(df4),left).drop(not required columns used in join)
join3_df.persist()
join4_df = join3_df.join(broadcast(df5),left).drop(not required columns used in join)
join4_df.persist()
join5_df = join4_df.join(broadcast(df6),left).drop(not required columns used in join)
join5_df.persist()
join6_df = join5_Df.join(broadcast(df7),left).drop(not required columns used in join)
join6_df.persist()
join7_df = join5_Df.join(broadcast(df8),left).drop(not required columns used in join)
join7_df.persist()
join8_df = join5_Df.join(broadcast(df9),left).drop(not required columns used in join)
join8_df.persist()
join9_df = join5_Df.join(broadcast(df10),left).drop(not required columns used in join)
join9_df.persist()
join10_df = join5_Df.join(broadcast(df11),left).drop(not required columns used in join)
join10_df.persist()
join11_df = join5_Df.join(broadcast(df12),left).drop(not required columns used in join)
join11_df.persist()
join12_df = join5_Df.join(broadcast(df13),left).drop(not required columns used in join)
join12_df.persist()

final_df = join5_Df.join(broadcast(df14),left).drop(not required columns used in join)
final_df.persist()

注:代码存在笔误,join5_Df应为join5_df,当前逻辑会重复基于join5_df连接后续表,导致计算冗余

最终写入情况

最终DataFrame含56列,以Parquet格式写入S3:

final_df.write.mode("overwrite").parquet(s3 location)
  • final_df大小:109GB
  • 写入耗时:7小时
  • S3生成文件:201个,单文件约570MB
  • 总记录数:8471166455

当前集群与Spark配置

集群信息

aws.emr.instance_type: "r5.12xlarge"
aws.emr.no_of_instances: "40"

Spark提交配置

spark-submit --deploy-mode client --master yarn 
--conf spark.shuffle.service.enabled=true 
--conf spark.dynamicAllocation.enabled=true 
--conf spark.dynamicAllocation.initialExecutors=5 
--conf spark.dynamicAllocation.minExecutors=5 
--conf spark.dynamicAllocation.maxExecutors=40 
--conf spark.executor.cores=10 
--conf spark.executor.memory=150G 
--conf spark.driver.memory=150G 
--conf "spark.driver.extraJavaOptions=-XX:+UseG1GC" 
--conf "spark.executor.extraJavaOptions=-XX:+UseG1GC" 

优化方案

一、代码逻辑优化

  1. 修正连接笔误:将所有join5_Df改为join5_df,基于前一次连接结果递进式连接,避免重复计算
  2. 精细化持久化策略:仅在连接完大表(如df4、df9)后进行持久化,且指定存储级别为MEMORY_AND_DISK_SER(比默认更节省内存与IO);若final_df仅用于写入,可直接去掉final_df.persist(),减少磁盘IO
  3. 调整连接顺序:优先连接小表(df12、df13、df7、df8),再连接中表、大表,减少中间结果的数据膨胀
  4. 移除不必要的distinct:若原始Parquet文件无重复数据,直接删除distinct()操作,降低计算开销

二、Spark配置调整

针对r5.12xlarge实例(48核、192G内存),优化后的提交配置:

spark-submit --deploy-mode cluster --master yarn 
--conf spark.shuffle.service.enabled=true 
--conf spark.dynamicAllocation.enabled=true 
--conf spark.dynamicAllocation.initialExecutors=10 
--conf spark.dynamicAllocation.minExecutors=10 
--conf spark.dynamicAllocation.maxExecutors=40 
--conf spark.executor.cores=5 
--conf spark.executor.memory=32G 
--conf spark.executor.memoryOverhead=4G 
--conf spark.driver.cores=8 
--conf spark.driver.memory=64G 
--conf spark.driver.memoryOverhead=8G 
--conf "spark.driver.extraJavaOptions=-XX:+UseG1GC -XX:MaxGCPauseMillis=200" 
--conf "spark.executor.extraJavaOptions=-XX:+UseG1GC -XX:MaxGCPauseMillis=200" 
--conf spark.sql.shuffle.partitions=2000 
--conf spark.sql.files.maxPartitionBytes=128m 
--conf spark.hadoop.mapreduce.output.fileoutputformat.compress=true 
--conf spark.hadoop.mapreduce.output.fileoutputformat.compress.codec=org.apache.hadoop.io.compress.SnappyCodec 
--conf spark.sql.parquet.compression.codec=snappy 
--conf spark.dynamicAllocation.executorIdleTimeout=60s 
--conf spark.shuffle.io.retryWait=5s 
--conf spark.shuffle.io.maxRetries=10 
--conf spark.hadoop.fs.s3a.fast.upload=true

配置说明:

  • 切换为cluster部署模式:避免client模式下本地驱动的网络瓶颈,适配大规模作业
  • 优化Executor资源分配:单实例分配5核/Executor(单实例可运行9个Executor,剩余3核留给系统),32G内存+4G overhead,平衡并行度与内存稳定性
  • 调整shuffle分区数:设置为2000,匹配40个Executor×5核的并行度,减少单个分区数据量,避免写入任务过载
  • 开启Snappy压缩:降低写入S3的数据量,减少IO耗时
  • S3上传优化:开启spark.hadoop.fs.s3a.fast.upload=true,利用多部分上传提升写入速度
  • GC与动态调优:控制GC停顿时间,调整动态分配参数减少资源启动延迟

三、写入S3优化

  1. 控制文件大小:通过spark.sql.files.maxPartitionBytes=128m将单文件大小控制在128MB左右,避免大文件写入瓶颈
  2. 分区写入:若数据有合适的分区键(如日期、地域),使用final_df.write.partitionBy("partition_col")按分区写入,并行写入多个分区提升速度,同时优化后续Athena查询效率

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 10:21:12