DBT中基于最新更新时间处理SQL UNION ALL重复记录的方法
解决DBT增量数据重叠时的去重问题
场景说明
每日增量数据存储路径格式如下:
./././day=26/<time>/<name>.parquet
数据管道重试后,单日目录下会生成多个时间子文件夹,导致同一ID的记录重复且可能存在不同值。原DBT语句通过UNION ALL拼接所有目录数据,但因数据重叠报错,需基于ets(更新时间字段)保留每个ID的最新记录。
原SQL语句
{% set directory_path = get_dynamic_path(table_name = 'table name') %} {% for file_path in directory_path %} SELECT *, date_add(MILLISECOND, ets, '1970-01-01') as ets_converted FROM PARQUET.{{ file_path }} {% if not loop.last %}UNION ALL{% endif %} {% endfor %}
修改后的去重SQL
{% set directory_path = get_dynamic_path(table_name = 'table name') %} WITH unioned_data AS ( {% for file_path in directory_path %} SELECT *, date_add(MILLISECOND, ets, '1970-01-01') as ets_converted FROM PARQUET.{{ file_path }} {% if not loop.last %}UNION ALL{% endif %} {% endfor %} ), ranked_data AS ( SELECT *, ROW_NUMBER() OVER (PARTITION BY id ORDER BY ets_converted DESC) AS rn FROM unioned_data ) SELECT -- 替换为实际业务字段,或用* EXCLUDE (rn)(如Snowflake支持) id, column1, column2, ets, ets_converted FROM ranked_data WHERE rn = 1
逻辑说明
- 数据拼接:通过
unioned_dataCTE保留原逻辑,将所有时间子文件夹的记录拼接在一起 - 排序打标:在
ranked_dataCTE中,用ROW_NUMBER()窗口函数按id分组,按转换后的更新时间ets_converted倒序排序,每个ID的最新记录会被标记为rn=1 - 筛选去重:最后只保留
rn=1的记录,得到每个ID的最新版本
注意事项
- 将
id替换为你的实际主键字段 - 若SQL方言支持(如Snowflake),可使用
* EXCLUDE (rn)代替手动列出字段,简化语句 - 确认
ets_converted的时间转换逻辑正确,确保排序顺序符合实际更新时间
内容的提问来源于stack exchange,提问作者Tushar Singh
相关产品推荐
相关产品推荐

