项目疑问:S3中CSV被删除时,AWS Glue能否同步删除对应Parquet数据?
实现S3 CSV删除同步清理对应Parquet数据的方案
针对你遇到的问题,核心要解决两个关键点:捕获S3文件删除事件和根据Parquet的存储策略处理数据清理,以下是具体实现步骤:
1. 捕获S3 CSV文件删除事件
首先给存储CSV的S3桶配置事件通知,触发后续的清理逻辑:
- 进入S3桶的「属性」→「事件通知」,创建新通知:
- 事件类型选择「删除」(包含永久删除和版本删除)
- 目标选择「Lambda函数」,创建或关联一个Lambda函数用于后续逻辑触发
- Lambda函数会收到S3事件的JSON payload,从中可以提取被删除的CSV文件的
key(即S3路径)
2. 根据Parquet存储策略处理清理
Parquet的不可变性意味着无法直接修改文件内容,只能根据你的数据生成方式选择不同的清理方案:
场景A:CSV与Parquet文件一对一映射
如果你的Glue作业是将每个CSV文件单独转换为对应的Parquet文件(比如保持相同文件名、仅替换后缀和存储路径),那么清理逻辑非常直接:
- 在Lambda函数中,根据被删CSV的路径生成对应的Parquet文件路径
- 调用S3 API直接删除该Parquet文件
示例Lambda代码(Python):
import boto3 s3_client = boto3.client('s3') PARQUET_BUCKET = "your-parquet-bucket-name" INPUT_PREFIX = "csv-input/" OUTPUT_PREFIX = "parquet-output/" def lambda_handler(event, context): # 提取被删除的CSV文件路径 deleted_csv_key = event['Records'][0]['s3']['object']['key'] # 转换为对应的Parquet路径 parquet_key = deleted_csv_key.replace(INPUT_PREFIX, OUTPUT_PREFIX).replace('.csv', '.parquet') # 删除Parquet文件 s3_client.delete_object(Bucket=PARQUET_BUCKET, Key=parquet_key)
场景B:多个CSV合并为Parquet文件
如果Glue作业是将多个CSV合并成少量Parquet文件(或分区表),则需要通过Glue作业重写过滤后的数据:
前置准备:初始转换时添加源文件标识
在最初的Glue转换作业中,必须给Parquet数据添加source_file字段,记录每条数据来自哪个CSV文件,这样后续才能精准过滤:
from awsglue.context import GlueContext from pyspark.context import SparkContext from pyspark.sql.functions import input_file_name sc = SparkContext() glue_context = GlueContext(sc) # 读取S3中的CSV数据,同时获取源文件名 df = glue_context.create_dynamic_frame.from_options( connection_type="s3", connection_options={"paths": ["s3://your-csv-bucket/csv-input/"], "recurse": True}, format="csv", format_options={"withHeader": True} ).toDF() # 添加source_file字段,存储原始CSV文件的S3路径 df = df.withColumn("source_file", input_file_name()) # 写入Parquet(如果是分区表,可添加partitionBy参数) df.write.parquet("s3://your-parquet-bucket/parquet-output/", mode="append")
触发清理流程
当CSV被删除时,Lambda函数将被删文件路径作为参数启动Glue作业,Glue作业执行以下步骤:
- 读取当前所有Parquet数据
- 过滤掉
source_file等于被删CSV路径的数据 - 用原子替换的方式重写Parquet数据(避免查询不一致)
示例Glue作业代码(Python):
import sys import boto3 from awsglue.utils import getResolvedOptions from awsglue.context import GlueContext from pyspark.context import SparkContext sc = SparkContext() glue_context = GlueContext(sc) s3_client = boto3.client('s3') PARQUET_BUCKET = "your-parquet-bucket-name" OUTPUT_PATH = "parquet-output/" TEMP_OUTPUT_PATH = "parquet-output-temp/" # 获取Lambda传递的被删文件参数 args = getResolvedOptions(sys.argv, ['DELETED_CSV_KEY']) deleted_csv_key = args['DELETED_CSV_KEY'] # 读取Parquet数据并过滤 df = glue_context.read.parquet(f"s3://{PARQUET_BUCKET}/{OUTPUT_PATH}") filtered_df = df.filter(df.source_file != f"s3://your-csv-bucket/{deleted_csv_key}") # 写入临时路径 filtered_df.write.parquet(f"s3://{PARQUET_BUCKET}/{TEMP_OUTPUT_PATH}", mode="overwrite") # 删除原路径所有文件 paginator = s3_client.get_paginator('list_objects_v2') for page in paginator.paginate(Bucket=PARQUET_BUCKET, Prefix=OUTPUT_PATH): if 'Contents' in page: for obj in page['Contents']: s3_client.delete_object(Bucket=PARQUET_BUCKET, Key=obj['Key']) # 将临时路径文件移动到原路径 for page in paginator.paginate(Bucket=PARQUET_BUCKET, Prefix=TEMP_OUTPUT_PATH): if 'Contents' in page: for obj in page['Contents']: new_key = obj['Key'].replace(TEMP_OUTPUT_PATH, OUTPUT_PATH) s3_client.copy_object( Bucket=PARQUET_BUCKET, Key=new_key, CopySource={'Bucket': PARQUET_BUCKET, 'Key': obj['Key']} ) s3_client.delete_object(Bucket=PARQUET_BUCKET, Key=obj['Key'])
3. Athena查询一致性保障
- 重写Parquet时必须用原子替换(先写临时路径,再替换原路径),避免Athena查询到部分更新的数据
- 如果使用分区表,仅需重写被删CSV所属的分区,无需全量处理,能大幅提升效率
内容的提问来源于stack exchange,提问作者Nicholas Gati
相关产品推荐
相关产品推荐

