AWS Glue作业后Parquet处理表出现重复记录求助
问题根源分析
你的重复问题大概率来自两个核心原因:
- ETL写入模式错误:Processed表采用了默认的
Append追加模式,每次ETL运行都会将处理后的数据追加到现有表中,而非替换或合并现有记录。即使你对当前批次数据做了去重,历史的旧记录依然留在Processed表中,最终导致重复。 - 原始表数据冗余:Glue爬虫默认是追加式更新元数据,如果你的S3文件是覆盖式更新,但爬虫没有配置为
更新现有表并删除旧分区/记录,原始表可能会保留同一TransportId的多条历史记录(比如旧的JSON文件记录没有被清理)。此时你的去重逻辑只处理了当前读取的数据集,但如果原始表本身就包含重复,写入Processed时追加就会累积重复。
另外,你提到的单独去重作业无效,很大概率也是因为该作业同样用了追加模式,没有真正替换Processed表的内容。
可行解决方案
方案1:修改ETL写入模式为覆盖(全量重写)
如果你的原始表是全量最新数据(即爬虫每次运行后,原始表只保留每个TransportId的最新记录),那么直接将Processed表的写入模式改为Overwrite,确保每次ETL运行都用最新的去重数据完全替换旧表:
# 在Glue ETL作业的写入节点添加配置 glueContext.write_dynamic_frame.from_catalog( frame=deduped_dynamic, database="your_db", table_name="Processed", additional_options={"writeBehavior": "overwrite"} )
注意:如果Processed表是分区表,可使用
partitionKeys指定分区,配合overwrite模式只覆盖对应分区,避免全表重写的性能损耗。
方案2:合并更新(增量+历史去重)
如果原始表是增量更新(保留历史记录),或者你不想全量重写Processed表,可采用合并更新的方式,基于TransportId和UnixTimeStamp保留最新记录:
- 读取Processed表的现有数据
- 读取当前批次的原始数据(已去重)
- 将两者合并后,再次按
TransportId分区、UnixTimeStamp降序取最新记录 - 用Overwrite模式写入Processed表
示例代码:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 读取现有Processed表数据 processed_df = glueContext.create_dynamic_frame.from_catalog( database="your_db", table_name="Processed" ).toDF() # 读取当前原始数据并去重 raw_df = AWSGlueDataCatalog_node1744383413520.toDF() window = Window.partitionBy("transportid").orderBy(F.desc("unixtimestamp")) raw_deduped_df = raw_df.withColumn("rn", F.row_number().over(window)).filter("rn == 1").drop("rn") # 合并数据并再次去重 combined_df = processed_df.unionByName(raw_deduped_df, allowMissingColumns=True) final_deduped_df = combined_df.withColumn("rn", F.row_number().over(window)).filter("rn == 1").drop("rn") # 转换为DynamicFrame并写入 final_dynamic = DynamicFrame.fromDF(final_deduped_df, glueContext, "final_deduped") glueContext.write_dynamic_frame.from_catalog( frame=final_dynamic, database="your_db", table_name="Processed", additional_options={"writeBehavior": "overwrite"} )
方案3:修正Glue爬虫配置
确保爬虫运行时会更新原始表的现有记录,而非追加:
- 在爬虫的
配置页面,找到高级选项,将当检测到现有分区中的更改时设置为更新表定义和数据 - 如果你的S3文件是覆盖式更新,可开启
删除不再存在于S3中的对象的分区,避免原始表保留无效的历史记录
方案4:使用Glue的Merge Into功能(适合大表增量更新)
对于数据量较大的场景,可使用Glue的Merge Into语法,直接在Processed表上执行合并更新,只更新有变化的TransportId记录:
MERGE INTO Processed p USING ( SELECT transportid, orgid, orgname, transportstatus, unixtimestamp FROM raw_table QUALIFY ROW_NUMBER() OVER (PARTITION BY transportid ORDER BY unixtimestamp DESC) = 1 ) r ON p.transportid = r.transportid WHEN MATCHED THEN UPDATE SET p.transportstatus = r.transportstatus, p.unixtimestamp = r.unixtimestamp, p.orgname = r.orgname WHEN NOT MATCHED THEN INSERT *
该方式需要Processed表支持ACID事务,需在表属性中开启
transactional=true。
验证步骤
- 先检查原始表是否存在同一
TransportId的多条记录:查询raw_table中transportid=654的记录数 - 检查ETL作业的写入配置,确认是否为
Append模式 - 运行一次全量去重+Overwrite写入的ETL作业,验证Processed表的重复是否消失
内容的提问来源于stack exchange,提问作者Max Manitskov
相关产品推荐
相关产品推荐

