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"
优化方案
一、代码逻辑优化
- 修正连接笔误:将所有
join5_Df改为join5_df,基于前一次连接结果递进式连接,避免重复计算 - 精细化持久化策略:仅在连接完大表(如df4、df9)后进行持久化,且指定存储级别为
MEMORY_AND_DISK_SER(比默认更节省内存与IO);若final_df仅用于写入,可直接去掉final_df.persist(),减少磁盘IO - 调整连接顺序:优先连接小表(df12、df13、df7、df8),再连接中表、大表,减少中间结果的数据膨胀
- 移除不必要的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优化
- 控制文件大小:通过
spark.sql.files.maxPartitionBytes=128m将单文件大小控制在128MB左右,避免大文件写入瓶颈 - 分区写入:若数据有合适的分区键(如日期、地域),使用
final_df.write.partitionBy("partition_col")按分区写入,并行写入多个分区提升速度,同时优化后续Athena查询效率
内容的提问来源于stack exchange,提问作者diksha yadgire
相关产品推荐
相关产品推荐

