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

项目疑问: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作业执行以下步骤:

  1. 读取当前所有Parquet数据
  2. 过滤掉source_file等于被删CSV路径的数据
  3. 用原子替换的方式重写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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 05:47:06