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

AWS Glue作业后Parquet处理表出现重复记录求助

问题根源分析

你的重复问题大概率来自两个核心原因:

  1. ETL写入模式错误:Processed表采用了默认的Append追加模式,每次ETL运行都会将处理后的数据追加到现有表中,而非替换或合并现有记录。即使你对当前批次数据做了去重,历史的旧记录依然留在Processed表中,最终导致重复。
  2. 原始表数据冗余: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保留最新记录:

  1. 读取Processed表的现有数据
  2. 读取当前批次的原始数据(已去重)
  3. 将两者合并后,再次按TransportId分区、UnixTimeStamp降序取最新记录
  4. 用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。

验证步骤
  1. 先检查原始表是否存在同一TransportId的多条记录:查询raw_table中transportid=654的记录数
  2. 检查ETL作业的写入配置,确认是否为Append模式
  3. 运行一次全量去重+Overwrite写入的ETL作业,验证Processed表的重复是否消失

内容的提问来源于stack exchange,提问作者Max Manitskov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 09:59:53