如何避免AWS Glue作业向Amazon Redshift数据库传输重复数据
避免AWS Glue向Redshift重复传输S3 CSV文件的解决方案
以下是几种实用的解决思路,覆盖不同场景:
1. 用S3事件触发作业,只处理新上传的文件
别让作业定时扫描整个桶,直接给S3桶配置事件通知(比如检测PutObject动作),新CSV一上传就触发Glue作业。作业里通过事件参数拿到新文件的S3路径,只加载这个路径的文件,从根源避免全量重复处理。
2. 维护已处理文件的元数据记录
建一个小型元数据表(可以存在Redshift、DynamoDB或者Glue数据目录里),每次处理完文件就把文件的S3路径或文件名存进去。作业启动时先查这个表,过滤掉已经处理过的文件,只处理未记录的。
示例Spark代码逻辑:
# 从Redshift的元数据表获取已处理文件列表 processed_files = spark.sql("SELECT s3_path FROM processed_files").rdd.map(lambda x: x[0]).collect() # 获取S3桶里的所有待处理文件路径(可通过boto3枚举实现get_all_s3_paths方法) all_files = get_all_s3_paths("s3://your-bucket/csv-files/") # 过滤出未处理的文件 unprocessed_files = [f for f in all_files if f not in processed_files] if unprocessed_files: # 加载未处理文件 df = spark.read.csv(unprocessed_files, header=True) # 写入Redshift df.write.mode("append").jdbc("jdbc:redshift://your-cluster-url:5439/db", "target_table", properties=your_db_props) # 将刚处理的文件路径写入元数据表 processed_df = spark.createDataFrame([(path,) for path in unprocessed_files], ["s3_path"]) processed_df.write.mode("append").jdbc("jdbc:redshift://your-cluster-url:5439/db", "processed_files", properties=your_db_props)
3. 依赖Redshift的主键/唯一键去重
如果CSV数据里有唯一标识(比如订单ID、日期+用户ID),给Redshift目标表加主键或唯一约束。Glue写入用append模式,之后在Redshift端处理重复:
- 写入时用
INSERT ... ON CONFLICT语法,冲突时跳过或更新:INSERT INTO target_table (col1, col2, unique_id) VALUES (...) ON CONFLICT (unique_id) DO NOTHING; -- 或者DO UPDATE SET ... - 或者定期跑去重SQL,保留最新的一条记录:
DELETE FROM target_table WHERE (unique_id, load_time) NOT IN ( SELECT unique_id, MAX(load_time) FROM target_table GROUP BY unique_id );
也可以在Glue脚本里先对DataFrame去重:
df = df.dropDuplicates(["unique_id"])
4. 按分区做增量加载
如果CSV是按日期分区存储的(比如s3://bucket/csv/date=2024-05-20/),作业只处理当天的分区路径。可以在脚本里动态生成当天的路径:
from datetime import datetime today = datetime.today().strftime("%Y-%m-%d") target_path = f"s3://your-bucket/csv/date={today}/" df = spark.read.csv(target_path, header=True)
这样每次作业只扫当天的文件,不会碰历史数据。
5. 处理完文件后移动或标记
处理完的文件直接移到S3的processed子目录,或者给文件名加个_processed后缀。作业运行时只扫描未处理目录的文件,处理完就执行移动操作:
import boto3 s3 = boto3.client("s3") bucket_name = "your-bucket" source_dir = "unprocessed-csv/" dest_dir = "processed-csv/" # 列出未处理目录的文件 response = s3.list_objects_v2(Bucket=bucket_name, Prefix=source_dir) for obj in response.get("Contents", []): file_key = obj["Key"] # 加载并处理文件 df = spark.read.csv(f"s3://{bucket_name}/{file_key}", header=True) df.write.mode("append").jdbc(...) # 移动文件到已处理目录 new_key = dest_dir + file_key.split('/')[-1] s3.copy_object(Bucket=bucket_name, CopySource=f"{bucket_name}/{file_key}", Key=new_key) s3.delete_object(Bucket=bucket_name, Key=file_key)
内容的提问来源于stack exchange,提问作者Srikanth
相关产品推荐
相关产品推荐

