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

如何解决AWS Glue ETL作业重复生成Parquet文件的问题

我之前在处理电商用户行为数据的Glue ETL任务时,刚好踩过这个文件累加的坑,给你分享两种可行的解决方案,都是基于Glue生成的Python脚本修改,完全符合你的需求:


1. 全量覆盖:直接替换所有现有Parquet文件

这是最直接的方案,适配你的场景——每次运行都想保留和源文件数量一致的Parquet文件(10个对应10个)。核心是把Glue默认的追加写入模式改成覆盖模式,这样每次运行都会先清空目标S3路径下的所有旧Parquet文件,再写入新的转换结果。

修改步骤:

找到你Glue脚本里的write_dynamic_frame.from_options代码块,添加format_options={"writeMode": "overwrite"}参数即可:

from awsglue.dynamicframe import DynamicFrame
from awsglue.context import GlueContext
from pyspark.context import SparkContext

sc = SparkContext.getOrCreate()
glueContext = GlueContext(sc)

# 读取源JSON/CSV文件(替换为你的源路径和格式)
source_df = glueContext.create_dynamic_frame.from_options(
    connection_type="s3",
    connection_options={"paths": ["s3://your-source-bucket/json-csv-files/"]},
    format="json"  # 或"csv",根据源文件类型调整
)

# 这里是你的转换逻辑(比如字段映射、清洗等)
transformed_df = source_df.apply_mapping([
    ("id", "string", "user_id", "string"),
    ("timestamp", "string", "event_time", "timestamp")
])

# 修改写入逻辑为覆盖模式
glueContext.write_dynamic_frame.from_options(
    frame=transformed_df,
    connection_type="s3",
    connection_options={"path": "s3://your-target-bucket/parquet-files/"},
    format="parquet",
    format_options={"writeMode": "overwrite"}  # 关键参数:开启覆盖
)

注意事项:

  • 这个模式会删除目标路径下的所有文件,请确保目标桶路径仅存放该任务的Parquet结果,不要混放其他数据。
  • 如果目标路径使用了分区(比如按日期分区),writeMode="overwrite"会覆盖整个分区路径;若需精准覆盖特定分区,可在connection_options中指定partitionKeys和partitionValues。

2. 增量处理:仅转换更新/新增的源文件

如果你的源文件不是每次都全部更新(比如只有部分JSON/CSV被修改或新增),这种方案更高效,能避免重复处理未变更的文件,同时只更新对应的Parquet文件。

方案A:使用Glue作业书签(官方推荐)

Glue内置的书签功能会自动记录已处理的文件,下次运行时仅处理新增或修改的文件。

配置步骤:

  1. 打开你的Glue作业配置,找到作业书签选项,设置为启用。
  2. 修改脚本的源文件读取逻辑,添加transformation_ctx标记(书签依赖此标记追踪状态):
# 读取源文件时添加transformation_ctx
source_df = glueContext.create_dynamic_frame.from_options(
    connection_type="s3",
    connection_options={
        "paths": ["s3://your-source-bucket/json-csv-files/"],
        "recurse": True  # 可选:若源文件在子目录下
    },
    format="json",
    transformation_ctx="source_df"  # 必须指定,书签核心依赖
)

# 转换逻辑...

# 写入时用追加模式(仅新增/更新对应Parquet文件)
glueContext.write_dynamic_frame.from_options(
    frame=transformed_df,
    connection_type="s3",
    connection_options={"path": "s3://your-target-bucket/parquet-files/"},
    format="parquet",
    format_options={"writeMode": "append"}
)

注意事项:

  • 书签与作业名称绑定,请勿随意修改作业名称或删除作业,否则会丢失追踪记录。
  • 若需重新处理所有文件,可在作业配置中点击重置书签。

方案B:自定义文件追踪(精准控制)

如果你需要更精细的控制(比如只覆盖被修改的源文件对应的Parquet),可以用boto3手动追踪文件的最后修改时间,筛选出需要处理的文件后,先删除旧Parquet再写入新文件:

import boto3
from datetime import datetime

# 初始化S3客户端
s3_client = boto3.client('s3')

# 读取上次运行时间(可存放在S3的标记文件或Glue参数中)
last_run_obj = s3_client.get_object(Bucket="your-target-bucket", Key="last_run_time.txt")
last_run_time = datetime.strptime(last_run_obj['Body'].read().decode(), "%Y-%m-%d %H:%M:%S")

# 筛选源桶中修改时间晚于上次运行时间的文件
response = s3_client.list_objects_v2(Bucket="your-source-bucket", Prefix="json-csv-files/")
filtered_files = []
for obj in response.get('Contents', []):
    if obj['LastModified'].replace(tzinfo=None) > last_run_time:
        filtered_files.append(f"s3://your-source-bucket/{obj['Key']}")

# 仅处理筛选后的文件
if filtered_files:
    source_df = glueContext.create_dynamic_frame.from_options(
        connection_type="s3",
        connection_options={"paths": filtered_files},
        format="json"
    )
    # 转换逻辑...
    
    # 删除对应的旧Parquet文件(假设文件名对应:xxx.json → xxx.parquet)
    for file_path in filtered_files:
        parquet_key = file_path.replace("json-csv-files/", "parquet-files/").replace(".json", ".parquet")
        s3_client.delete_object(Bucket="your-target-bucket", Key=parquet_key.split("s3://your-target-bucket/")[1])
    
    # 写入新的Parquet文件
    glueContext.write_dynamic_frame.from_options(
        frame=transformed_df,
        connection_type="s3",
        connection_options={"path": "s3://your-target-bucket/parquet-files/"},
        format="parquet",
        format_options={"writeMode": "append"}
    )

# 更新上次运行时间到S3标记文件
current_time = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
s3_client.put_object(Bucket="your-target-bucket", Key="last_run_time.txt", Body=current_time)

内容的提问来源于stack exchange,提问作者Sparta

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 04:18:25