如何解决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内置的书签功能会自动记录已处理的文件,下次运行时仅处理新增或修改的文件。
配置步骤:
- 打开你的Glue作业配置,找到作业书签选项,设置为启用。
- 修改脚本的源文件读取逻辑,添加
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
相关产品推荐
相关产品推荐

