AWS Glue作业节点配置与OOM问题排查:70文件夹合并任务异常
问题分析与解决建议
问题根源分析
- Driver端线程池与Spark调度冲突:Spark本身是分布式调度框架,手动使用
concurrent.futures.ThreadPoolExecutor在Driver端开启线程并行提交任务,会导致Driver内存被大量并发的任务元数据处理、请求调度占满,同时打乱Spark原生的资源调度逻辑,导致Executor资源无法有效利用(监控显示活跃Executor仅1.5就是证明),最终触发OOM。 - 单个Task内存过载:每个文件夹700个文件,
repartition(1)会将该文件夹下所有数据shuffle到单个Executor的单个Task中处理,若单文件夹数据量较大,单个Task的内存占用会远超Executor的内存阈值;当线程池同时触发多个这类Task时,Executor内存会被直接打满,而处理到第50个文件夹时,可能剩余文件夹的平均文件更大,或前面任务未及时释放内存,最终触发异常。 - 资源配置未匹配任务需求:G2X节点虽有64GB内存,但默认的Executor/Driver内存配置可能未针对大Task场景调优,导致内存分配不足;同时最大工作节点4、最大并发3的配置未被有效利用,资源闲置与过载并存。
具体解决建议
1. 移除线程池,改用Spark原生分布式处理
删除ThreadPoolExecutor相关代码,将文件夹列表转换为RDD/DataFrame的分区,让Spark原生调度Executor并行处理:
# 示例:获取所有文件夹路径,转为RDD后分布式处理 folder_paths = get_all_folder_paths() # 自定义方法获取所有目标文件夹 sc.parallelize(folder_paths, numSlices=10).foreach(process_folder) def process_folder(folder_path): df = spark.read.parquet(folder_path) # 可选:过滤/清洗数据减少内存占用 df.coalesce(1).write.parquet(f"{target_path}/{folder_path.split('/')[-1]}")
这种方式让Driver仅负责调度,具体任务由Executor分布式执行,避免Driver内存过载,同时充分利用集群资源。
2. 优化数据合并逻辑
- 用
coalesce(1)替代repartition(1):coalesce无需shuffle,直接合并分区,大幅减少内存开销;仅当数据分布极不均匀时才考虑repartition。 - 先过滤无效数据:读取文件后先执行过滤、去重、列裁剪等操作,减少需要合并的数据量,降低内存压力。
3. 调优Glue作业资源配置
- Driver端配置:在作业参数中添加:
增大Driver内存,避免调度任务时内存不足。--driver-memory 32g --driver-cores 4 - Executor端配置:针对G2X节点(64GB内存),设置:
给Executor预留足够内存处理大Task,同时匹配最大并发3的配置。--executor-memory 48g --executor-cores 12 --num-executors 3 - Spark内存管理调优:添加参数:
分配更多内存给执行计算,减少内存溢出风险。--conf spark.memory.fraction=0.7 --conf spark.memory.storageFraction=0.5
4. 分批处理文件夹
若分布式处理仍有问题,将70个文件夹拆分为多批(比如每批10个),每批处理完成后再启动下一批,避免同时运行过多大内存任务,让集群有时间释放中间资源。
5. 精准排查异常
- 查看Glue作业日志,确认OOM发生在Driver还是Executor:
- 若为Driver:重点优化Driver内存和移除线程池;
- 若为Executor:重点调优Executor内存和数据处理逻辑。
- 检查剩余20个文件夹的文件大小,确认是否存在超大文件,针对性单独处理(比如拆分后再合并)。
内容的提问来源于stack exchange,提问作者Rimpa Dey Sarkar
相关产品推荐
相关产品推荐

