DBT Snapshot处理仅追加源表非唯一记录的方案咨询
DBT 仅追加源表多变更记录快照实现方案
仅追加模式数据湖的源表单次同步周期内同一主键存在多条变更的场景非常普遍,DBT 原生 Snapshot 默认单主键仅处理最新一条变更的逻辑无法匹配需求,可通过以下两种方案实现预期效果:
方案1:预加工源数据 + DBT 原生Snapshot(轻量适配)
先通过中间层预处理源表,提前按主键排序计算每条变更的失效时间,再喂给原生 Snapshot 处理:
-- 预处理模型 stg_source_records_ranked.sql {{ config(materialized='view') }} with source_all as ( select id, some_attribute, updated_at from {{ source('your_schema', 'your_source_table') }} -- 过滤已加载过的旧记录提升性能 where updated_at > (select coalesce(max(dbt_valid_from), '1970-01-01'::timestamp) from {{ ref('your_snapshot') }}) ) select id, some_attribute, updated_at, -- 取同主键下一条记录的更新时间作为当前记录的失效时间 lead(updated_at) over(partition by id order by updated_at asc) as next_updated_at from source_all
再配置 Snapshot,使用 timestamp 策略:
{% snapshot your_snapshot %} {{ config( target_schema='your_snapshot_schema', unique_key='id', strategy='timestamp', updated_at='updated_at' ) }} select id, some_attribute, updated_at from {{ ref('stg_source_records_ranked') }} -- 仅保留每个批次最新的记录给原生Snapshot更新旧的活跃记录 where next_updated_at is null {% endsnapshot %}
处理完成后再跑一个下游模型补全中间的历史变更记录即可。
方案2:增量表自定义快照逻辑(更推荐,完全匹配需求)
直接用增量表实现完整快照逻辑,无需依赖DBT原生Snapshot的内置规则,完全适配多变更场景:
-- 增量快照模型 your_custom_snapshot.sql {{ config( materialized='incremental', incremental_strategy='merge', unique_key=['id', 'dbt_valid_from'], merge_update_columns=['dbt_valid_to'] ) }} -- 拉取本次需要处理的新增源记录 with new_records as ( select id, some_attribute, updated_at from {{ source('your_schema', 'your_source_table') }} {% if is_incremental() %} where updated_at > (select coalesce(max(dbt_valid_from), '1970-01-01'::timestamp) from {{ this }}) {% endif %} ), -- 计算新记录的有效期 ranked_new as ( select id, some_attribute, updated_at as dbt_valid_from, lead(updated_at) over(partition by id order by updated_at asc) as dbt_valid_to from new_records ), -- 拉取快照表中需要更新失效时间的原有活跃记录 active_old_records as ( select id, some_attribute, dbt_valid_from from {{ this }} where dbt_valid_to is null {% if is_incremental() %} and id in (select distinct id from new_records) {% endif %} ), -- 构造需要更新的旧活跃记录,将失效时间设为对应主键新记录的最小更新时间 updated_old as ( select a.id, a.some_attribute, a.dbt_valid_from, min(n.dbt_valid_from) as dbt_valid_to from active_old_records a join ranked_new n on a.id = n.id group by 1,2,3 ) -- 合并输出:要更新的旧记录 + 要插入的新变更记录 select * from updated_old union all select * from ranked_new
注意事项
- 若使用的数仓不支持merge增量策略,可调整为按ID或日期分区的append+分区覆写逻辑,核心计算规则无需修改
- 自定义增量方案比修改原生Snapshot宏的维护成本低,且可灵活适配各类特殊业务规则
内容的提问来源于stack exchange,提问作者Jose Bagatelli
相关产品推荐
相关产品推荐

