EMR环境下重跑PySpark作业为何耗时增加、任务数暴涨?
EMR PySpark二次运行任务数暴涨、性能下降根因
该现象是S3对象存储与Spark写入逻辑适配问题导致,与Spark内存缓存、Kernel本地状态无关,spark.catalog.clearCache()、重启Kernel操作不生效符合预期,具体触发原因如下:
- 覆盖写入阶段的前置目录扫描直接拉高任务数
首次运行作业时,目标S3写入路径为空,Spark执行覆盖写入逻辑时不需要做前置文件扫描、删除标记生成,任务规划阶段仅需根据上游数据的shuffle分区数(默认值为800)生成写入任务,因此总任务数稳定在800,运行耗时短。
第二次运行相同作业时,目标S3路径已存在首次写入的全量文件,Spark默认写入逻辑会先递归遍历目标路径下的所有对象、生成待删除文件列表,这一步遍历得到的S3对象会被纳入输入切片计算逻辑,直接推高总任务数。如果首次写入时未做小文件合并,路径下存在大量碎文件,切片数会直接涨到数千量级,观测到的5000任务规模完全匹配该场景特征。 - EMR默认元数据同步逻辑增加额外规划开销
EMR Notebook默认绑定Glue Data Catalog作为元存储,若写入目标是分区表,第二次运行作业时,Spark会先从Glue拉取目标表的所有历史分区元数据做校验,这部分元数据拉取、分区校验操作不仅会增加任务规划阶段耗时,拉取到的分区统计信息还会干扰Spark动态分区裁剪判断,导致*AQE(自适应查询执行)*无法正常合并小分区,最终任务数无法回落至正常水平。spark.catalog.clearCache()仅能清理Spark进程内存中缓存的表数据,无法清理Glue侧存储的表元数据、S3侧存储的实际对象信息;重启Kernel仅会重置本地Spark会话状态,不会修改外部存储的S3对象、Glue元数据,因此两类操作都无法解决问题。 - 默认S3提交器的冗余操作拉长运行时长
若未显式配置EMR优化的S3提交器,Spark会使用原生FileOutputCommitter,在S3对象存储上会产生大量临时文件重命名、目录列表操作,第二次写入时需要先清理上一次作业遗留的临时文件,这部分清理操作会生成额外任务,进一步拉长总运行时长。
验证与修复方案
- 验证方式:写入前手动清空目标S3路径后重跑作业,若任务数回落至800、耗时回到10秒级别,即可确认是上述原因导致。
- 修复配置:写入作业前添加以下参数,跳过不必要的全路径扫描、使用S3优化提交器降低冗余开销:
# 开启动态分区覆盖,避免全表扫描 spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic") # 启用EMR S3优化提交器 spark.conf.set("spark.hadoop.fs.s3a.committer.name", "magic") spark.conf.set("spark.hadoop.fs.s3a.committer.magic.enabled", "true") # 关闭非必要的提交阶段清理逻辑 spark.conf.set("spark.hadoop.mapreduce.fileoutputcommitter.cleanup.skipped", "true") # 控制写入文件大小,避免产生过多小文件 spark.conf.set("spark.sql.files.maxRecordsPerFile", 5000000)
- 分区表写入场景下,提前通过
where条件指定本次写入的分区范围,避免Spark自动扫描全表所有分区做校验。
内容的提问来源于stack exchange,提问作者EnterPassword
相关产品推荐
相关产品推荐

