如何在dbt snapshot中仅捕获指定列的历史变更数据?
实现差异化变更记录的dbt快照方案
可以实现,但无法通过dbt默认的check快照策略直接达成,需要自定义快照的merge逻辑,区分不同列的变更行为。
核心思路
默认的check策略会对所有指定列的变更生成新的快照行,导致频繁变更的column2快速膨胀数据量。要实现需求,需在快照的merge阶段添加判断逻辑:
- 当column3变更时:保留默认快照行为,生成新行并标记旧行过期。
- 当仅column2变更时:直接更新当前活跃行的column2值,不生成新的历史记录。
具体实现代码
以下是自定义快照的示例代码,需根据你的数据源和表结构调整:
{% snapshot custom_column_snapshot %} {{ config( target_schema='snapshots', unique_key='column1', strategy='check', check_cols=['column2', 'column3'], ) }} {% set source_relation = source('your_source_schema', 'your_source_table') %} {% set snapshot_relation = this %} -- 检测活跃行的变更类型 WITH active_snapshots AS ( SELECT * FROM {{ snapshot_relation }} WHERE dbt_valid_to IS NULL ), change_detection AS ( SELECT s.column1, CASE WHEN s.column3 != sd.column3 THEN 'column3_change' WHEN s.column2 != sd.column2 THEN 'column2_change' ELSE 'no_change' END AS change_type, sd.* FROM active_snapshots s JOIN {{ source_relation }} sd ON s.column1 = sd.column1 ), -- 准备column3变更时的插入数据 column3_change_rows AS ( SELECT *, CURRENT_TIMESTAMP AS dbt_valid_from, NULL AS dbt_valid_to, {{ dbt_utils.generate_surrogate_key(['column1', 'CURRENT_TIMESTAMP']) }} AS dbt_scd_id FROM change_detection WHERE change_type = 'column3_change' ), -- 准备column2变更时的更新数据(复用原有scd_id和生效时间) column2_change_rows AS ( SELECT sd.column1, sd.column2, sd.column3, s.dbt_scd_id, s.dbt_valid_from, NULL AS dbt_valid_to FROM {{ source_relation }} sd JOIN active_snapshots s ON sd.column1 = s.column1 JOIN change_detection cd ON sd.column1 = cd.column1 WHERE cd.change_type = 'column2_change' ) -- 执行merge操作 MERGE INTO {{ snapshot_relation }} t USING ( SELECT * FROM column3_change_rows UNION ALL SELECT * FROM column2_change_rows ) s ON t.dbt_scd_id = s.dbt_scd_id -- column3变更时,标记旧行过期 WHEN MATCHED AND t.dbt_valid_to IS NULL AND s.change_type = 'column3_change' THEN UPDATE SET t.dbt_valid_to = CURRENT_TIMESTAMP -- column3变更时插入新行,column2变更时更新现有行 WHEN NOT MATCHED THEN INSERT (column1, column2, column3, dbt_scd_id, dbt_valid_from, dbt_valid_to) VALUES (s.column1, s.column2, s.column3, s.dbt_scd_id, s.dbt_valid_from, s.dbt_valid_to) {% endsnapshot %}
关键逻辑说明
- 变更类型检测:仅针对快照表中未过期的活跃行,判断是column2还是column3发生变更。
- column3变更处理:生成新的快照行,同时将原活跃行的
dbt_valid_to设为当前时间,标记为过期,保留完整历史。 - column2变更处理:直接复用原活跃行的
dbt_scd_id和dbt_valid_from,更新column2的值,不生成新行,避免数据量激增。 - 每日运行适配:因为只处理活跃行,每日快照运行时不会重复处理已过期的历史数据,保证效率。
注意事项
- 需确保
dbt_utils包已安装(用于生成唯一的dbt_scd_id),若未安装可手动生成唯一键。 - 测试时需分别验证仅column2变更、仅column3变更、两列同时变更的场景,确保符合预期。
- 根据实际表结构调整插入和更新的字段列表,避免遗漏必要字段。
内容的提问来源于stack exchange,提问作者jvr
相关产品推荐
相关产品推荐

