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

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

逻辑说明

  1. 数据拼接:通过unioned_data CTE保留原逻辑,将所有时间子文件夹的记录拼接在一起
  2. 排序打标:在ranked_data CTE中,用ROW_NUMBER()窗口函数按id分组,按转换后的更新时间ets_converted倒序排序,每个ID的最新记录会被标记为rn=1
  3. 筛选去重:最后只保留rn=1的记录,得到每个ID的最新版本

注意事项

  • 将id替换为你的实际主键字段
  • 若SQL方言支持(如Snowflake),可使用* EXCLUDE (rn)代替手动列出字段,简化语句
  • 确认ets_converted的时间转换逻辑正确,确保排序顺序符合实际更新时间

内容的提问来源于stack exchange,提问作者Tushar Singh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 01:11:02