使用Glue做每日ETL增量加载时curated区数据翻倍问题如何解决
针对Glue增量ETL链路数据重复问题的整改方案
一、修复Raw到Curated区Upsert逻辑(解决重复根因)
- 明确主键与排序规则:无论采用Delta Lake、Iceberg等开放表格式还是普通分区表做Curated层存储,都必须指定全局唯一的业务主键(如订单ID、用户唯一标识),同时以记录最后更新时间戳为排序依据,合并时永远取最新版本覆盖旧数据,避免同主键多版本同时留存。若用Spark SQL实现upsert逻辑,优先使用标准
MERGE INTO语法,避免自定义join逻辑遗漏匹配条件导致旧数据未被覆盖。 - 收敛扫描分区范围:每次任务运行仅读取当前执行日期对应的Raw区增量数据(如1月11日运行仅加载
2020/01/11/路径下的文件),仅关联Curated区中与当前增量主键匹配的历史数据做合并,不要全量扫描所有历史分区,避免旧数据被重复写入。 - 前置单批次去重逻辑:写入Curated区前,先对当前批次的增量数据按主键去重,避免单批次内同主键存在多条记录导致合并异常。
二、调整Curated到Consumption区的增量输出规则
- 新增增量标识字段:在Curated区表结构中增加
dt(数据所属业务日期)、op_type(操作类型:I=新增/U=更新/D=删除)、load_ts(ETL任务加载时间戳)三个辅助字段,为消费层筛选增量提供明确依据。 - 固定增量筛选逻辑:Consumption层每次拉取数据时,仅筛选
load_ts大于上一次消费最大时间戳的记录,可额外增加按主键+load_ts倒序取第一条的逻辑,兜底过滤残留的重复数据。 - 若Curated区采用Delta Lake/Iceberg存储,可直接调用自带的变更数据捕获(CDC)能力,直接获取两次版本之间的delta记录,无需自行开发筛选逻辑。
三、清洗现有存量重复数据
执行一次性离线清洗任务:对Curated区现有数据按主键分组,仅保留load_ts最新的一条记录,其余重复数据删除,将总数据量恢复至预期的1100条。
内容的提问来源于stack exchange,提问作者Suganya
相关产品推荐
相关产品推荐

