AWS Glue拆分CSV Gzip文件时任务仅分配给单个worker无法分发
核心问题原因
gzip格式本身是不可分割的压缩类型,Glue底层基于Spark运行,单个gzip压缩文件只能由1个task完整读取、无法拆分到多个task并行处理,是导致你作业并行度低的核心原因:
- 你当前有4个5GB大小的gzip输入文件,读阶段默认最多生成4个处理task,如果你配置的worker数超过4,多余的worker天然就没有任务可以处理
- 你读阶段配置的
"groupFiles": "inPartition"是面向小文件合并场景的配置,会尝试把同一分区下的多个文件分配给同一个task处理,进一步压低了并行度 - 单个5GB的gzip文件解压后体积通常在20GB以上,如果你用的是默认G.1X规格(16GB内存)的worker,处理时会触发OOM异常,task会重试到其他可用worker,最终大概率所有task都挤到少量能承载负载的worker上运行,出现其余worker完全闲置的情况
排查&解决步骤
- 先调整作业worker配置:将worker规格升级为G.2X(32GB内存)及以上,worker数量配置为4~10个,保证单个worker可以承载单个5GB gzip文件的解压、处理负载,避免task频繁重试
- 修改读阶段配置:删除
create_dynamic_frame.from_catalog的additional_options中的"groupFiles": "inPartition"配置,或者改为"groupFiles": "none",避免Glue把多个大文件绑定到同一个task处理,保证4个输入文件可以对应4个独立的处理task - 手动重分区提升后续处理并行度:读入并完成字段映射之后,对动态帧执行重分区操作,把数据打散到更多task,让后续写出阶段可以用到所有worker,修改后的代码示例如下:
applymapping1 = ApplyMapping.apply(frame = datasource0, mappings = [ ("tagids", "string", "internal_tagids", "string"), ("channel", "long", "internal_channel", "long")], transformation_ctx = "applymapping1") # 新增重分区逻辑,并行度可以按worker数*2~3倍设置,比如你配了10个worker就设为20~30 repartitioned_frame = applymapping1.repartition(20)
- 修改写出配置:将
write_dynamic_frame.from_options的输入frame改为重分区后的frame,同时删除写出配置中的"groupFiles": "inPartition",保留"groupSize": "1073741824"即可,这样多个task可以并行写出1GB大小的目标文件 - 长期优化:如果可以调整输入文件的生成规则,建议把上游生成的5GB gzip文件拆分为1GB大小的小gzip文件,读阶段的天然并行度就会提升到20个task,不需要重分区即可跑满所有worker,避免shuffle开销
内容的提问来源于stack exchange,提问作者user1888955
相关产品推荐
相关产品推荐

