如何将AWS Glue作业输出文件重命名为.json/.parquet格式
解决AWS Glue作业输出指定名称文件的问题
我使用以下AWS Glue作业代码将数据写入AWS S3存储位置,但最终生成的是part文件。我的需求是将输出文件保存为指定名称的.json或.parquet格式,希望得到帮助。
用户提供的原始代码:
s3_loc = "s3a://s3_location/path" ##this is the default part of the glue script job = Job(glueContext) job.init(args['JOB_NAME'], args) datasource0 = glueContext.create_dynamic_frame.from_catalog(database = "dbname", table_name = "tableName", transformation_ctx = "datasource0") applymapping1 = ApplyMapping.apply(frame = datasource0, mappings = [("user id", "string", "user id", "string"), ("e-mail", "string", "e-mail", "string"), ("e-mail 2", "string", "e-mail 2", "string")], transformation_ctx = "applymapping1") timestampedDf = applymapping1.toDF().withColumn("export_timestamp", current_timestamp()) df = timestampedDf.withColumn("client", lit("122")).withColumn("partition_date", current_date()).coalesce(1) dyf = DynamicFrame.fromDF(df, glueContext, "dyf") print("---------> Let's go <---------") datasink2 = glueContext.write_dynamic_frame.from_options(frame = dyf, connection_type = "s3", connection_options = {"path": s3_loc , "partitionKeys": ["client","partition_date"]}, format = "json", transformation_ctx = "datasink2") print("---------> Let's finish <---------") job.commit()
核心原因
AWS Glue的write_dynamic_frame默认遵循Spark分布式输出规则,生成part-xxx格式的文件。要生成指定名称的文件,需额外添加重命名逻辑,同时要注意代码中配置的**分区键(client、partition_date)**会让文件落在对应分区目录下。
解决方案:Spark DataFrame API + S3重命名
直接使用Spark DataFrame的写入API配合S3操作,实现指定文件名输出:
from pyspark.sql.functions import current_timestamp, lit, current_date import boto3 from datetime import datetime s3_loc = "s3a://s3_location/path" target_filename = "user_data.json" # 自定义目标文件名,格式可改为.parquet # 初始化Glue作业 job = Job(glueContext) job.init(args['JOB_NAME'], args) # 原有数据处理逻辑保持不变 datasource0 = glueContext.create_dynamic_frame.from_catalog(database = "dbname", table_name = "tableName", transformation_ctx = "datasource0") applymapping1 = ApplyMapping.apply(frame = datasource0, mappings = [("user id", "string", "user id", "string"), ("e-mail", "string", "e-mail", "string"), ("e-mail 2", "string", "e-mail 2", "string")], transformation_ctx = "applymapping1") timestampedDf = applymapping1.toDF().withColumn("export_timestamp", current_timestamp()) df = timestampedDf.withColumn("client", lit("122")).withColumn("partition_date", current_date()).coalesce(1) # 步骤1:将数据写入临时路径 temp_path = f"{s3_loc}/temp_output" df.write.mode("overwrite").json(temp_path) # 若用parquet,改为.parquet(temp_path) # 步骤2:解析S3路径,准备重命名 s3_client = boto3.client('s3') bucket_name = s3_loc.split("//")[1].split("/")[0] base_prefix = "/".join(s3_loc.split("//")[1].split("/")[1:]) temp_prefix = f"{base_prefix}/temp_output" # 步骤3:定位临时路径下的part文件,重命名到目标分区目录 partition_date_str = datetime.now().strftime('%Y-%m-%d') target_key = f"{base_prefix}/client=122/partition_date={partition_date_str}/{target_filename}" response = s3_client.list_objects_v2(Bucket=bucket_name, Prefix=temp_prefix) for obj in response.get('Contents', []): if obj['Key'].endswith('.json') and 'part-' in obj['Key']: # parquet格式改为.endswith('.parquet') # 复制并重命名文件 s3_client.copy_object( Bucket=bucket_name, CopySource={'Bucket': bucket_name, 'Key': obj['Key']}, Key=target_key ) # 删除原part文件 s3_client.delete_object(Bucket=bucket_name, Key=obj['Key']) # 清理临时目录的_SUCCESS标记 s3_client.delete_object(Bucket=bucket_name, Key=f"{temp_prefix}/_SUCCESS") job.commit()
关键注意事项
coalesce(1)会将所有数据合并到一个分区,仅适合小数据量场景;大数据量下使用会导致性能瓶颈,建议保留多part文件或按需调整分区数。- 确保Glue作业的IAM角色拥有S3的
ListObjects、CopyObject、DeleteObject权限。 - 分区目录格式为
client=xxx/partition_date=yyyy-mm-dd,是Spark默认分区格式,若需自定义可调整写入路径逻辑。
内容的提问来源于stack exchange,提问作者tejaswinipatil
相关产品推荐
相关产品推荐

