You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.12 17:26:05