AWS Glue Python脚本无法生成指定大小文件问题
解决AWS Glue生成小文件问题(目标100-250MB gz文件)
问题背景
当前的AWS Glue Python脚本在S3桶中默认生成多个4MB的gzip压缩文件,尝试通过groupFiles和groupSize参数调整文件大小至100-250MB,但未生效。
当前脚本
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job import datetime args = getResolvedOptions(sys.argv, ['target_BucketName', 'JOB_NAME']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) outputbucketname = args['target_BucketName'] timestamp = datetime.datetime.now().strftime("%Y%m%d") filename = f"tbd{timestamp}" output_path = f"{outputbucketname}/{filename}" # Script generated for node AWS Glue Data Catalog AWSGlueDataCatalog_node075257312 = glueContext.create_dynamic_frame.from_catalog(database="ardt", table_name="_ard_tbd", transformation_ctx="AWSGlueDataCatalog_node075257312") # Script generated for node Amazon S3 AmazonS3_node075284688 = glueContext.write_dynamic_frame.from_options(frame=AWSGlueDataCatalog_node075257312, connection_type="s3", format="csv", format_options={"separator": "|"}, connection_options={"path": output_path, "compression": "gzip", "recurse": True, "groupFiles": "inPartition", "groupSize": "100000000"}, transformation_ctx="AmazonS3_node075284688") job.commit()
问题原因及解决方案
原因分析
groupFiles设为"inPartition"仅合并同分区内的小文件,若数据源本身分区粒度小,该参数无法生效;同时Spark默认分区数量过多,直接导致输出大量小文件。
具体解决步骤
重分区控制文件数量
将DynamicFrame转为Spark DataFrame后,根据目标文件大小计算合适的分区数(例如总数据压缩后为10GB,目标单文件200MB则设50个分区),再转回DynamicFrame:# 转换为DataFrame并重分区 df = AWSGlueDataCatalog_node075257312.toDF() df_repartitioned = df.repartition(50) # 可根据实际数据量调整分区数 dynamic_frame_repartitioned = DynamicFrame.fromDF(df_repartitioned, glueContext, "dynamic_frame_repartitioned")修正S3写入参数
- 将
groupFiles改为"merge",实现跨分区文件合并 - 调整
groupSize为目标大小(100MB=104857600字节,250MB=262144000字节) - 移除无用的
recurse参数(该参数用于读取而非写入)
- 将
调整作业资源配置
在Glue作业配置中,适当提高DPU数量和最大并发数,确保资源足够支持文件合并操作。
完整修改后的脚本
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job from awsglue.dynamicframe import DynamicFrame import datetime args = getResolvedOptions(sys.argv, ['target_BucketName', 'JOB_NAME']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) outputbucketname = args['target_BucketName'] timestamp = datetime.datetime.now().strftime("%Y%m%d") filename = f"tbd{timestamp}" output_path = f"{outputbucketname}/{filename}" # 读取数据源 AWSGlueDataCatalog_node075257312 = glueContext.create_dynamic_frame.from_catalog( database="ardt", table_name="_ard_tbd", transformation_ctx="AWSGlueDataCatalog_node075257312" ) # 重分区调整文件大小 df = AWSGlueDataCatalog_node075257312.toDF() df_repartitioned = df.repartition(50) # 根据实际数据量调整分区数 dynamic_frame_repartitioned = DynamicFrame.fromDF(df_repartitioned, glueContext, "dynamic_frame_repartitioned") # 写入S3并合并文件 AmazonS3_node075284688 = glueContext.write_dynamic_frame.from_options( frame=dynamic_frame_repartitioned, connection_type="s3", format="csv", format_options={"separator": "|"}, connection_options={ "path": output_path, "compression": "gzip", "groupFiles": "merge", "groupSize": "209715200" # 200MB,可按需调整为104857600或262144000 }, transformation_ctx="AmazonS3_node075284688" ) job.commit()
内容的提问来源于stack exchange,提问作者Marcus
相关产品推荐
相关产品推荐

