基于DBT/Snowflake实现模型Merge操作的行变更日志记录方案问询
DBT Merge操作后统计插入/更新行数并写入日志表的实现方案
完全可以通过DBT的Post Hook结合数据库原生功能实现需求,不需要额外工具包。以下是具体实现步骤:
1. 提前创建日志表
如果日志表log_table尚未存在,先执行DDL创建(以Snowflake为例,其他数据库语法类似):
CREATE TABLE log_table ( TABLE_LOADED VARCHAR(100) PRIMARY KEY, INSERTED_ROWS INTEGER, UPDATED_ROWS INTEGER, LOADED_AT TIMESTAMP DEFAULT CURRENT_TIMESTAMP() );
2. 编写Post Hook逻辑
核心思路是:在Merge操作完成后,通过数据库的系统函数/结果扫描获取本次Merge的插入、更新行数,再将数据写入日志表。不同数据库的实现细节略有差异,以下是主流数据库的示例:
示例1:Snowflake环境
Snowflake支持通过RESULT_SCAN扫描最后一次查询的执行结果,提取Merge操作的插入、更新行数。
步骤1:创建存储过程封装日志逻辑
CREATE OR REPLACE PROCEDURE log_merge_results(table_name VARCHAR) RETURNS VARCHAR LANGUAGE SQL AS $$ DECLARE insert_count INTEGER; update_count INTEGER; BEGIN -- 从最后一次Merge的执行结果中提取行数 SELECT "number of rows inserted", "number of rows updated" INTO insert_count, update_count FROM TABLE(RESULT_SCAN(LAST_QUERY_ID())); -- 合并更新日志表(存在则更新,不存在则插入) MERGE INTO log_table lt USING (SELECT table_name AS TABLE_LOADED, insert_count AS INSERTED_ROWS, update_count AS UPDATED_ROWS) src ON lt.TABLE_LOADED = src.TABLE_LOADED WHEN MATCHED THEN UPDATE SET INSERTED_ROWS = src.INSERTED_ROWS, UPDATED_ROWS = src.UPDATED_ROWS, LOADED_AT = CURRENT_TIMESTAMP() WHEN NOT MATCHED THEN INSERT (TABLE_LOADED, INSERTED_ROWS, UPDATED_ROWS) VALUES (src.TABLE_LOADED, src.INSERTED_ROWS, src.UPDATED_ROWS); RETURN '已记录表 ' || table_name || ':插入' || insert_count || '行,更新' || update_count || '行'; END; $$;
步骤2:在DBT模型中配置Post Hook
在需要执行Merge的模型配置中添加Post Hook,调用上述存储过程:
models: your_project: your_merge_model: materialized: incremental incremental_strategy: merge post-hook: - "CALL log_merge_results('{{ this.name }}');"
示例2:PostgreSQL环境
PostgreSQL通过INSERT ... ON CONFLICT实现Merge逻辑,可通过GET DIAGNOSTICS获取受影响行数,结合自定义变量区分插入、更新:
步骤1:编写Post Hook SQL
直接在模型的Post Hook中执行以下逻辑(无需存储过程):
{% set table_name = this.name %} DO $$ DECLARE insert_rows INTEGER; update_rows INTEGER; BEGIN -- 获取本次操作的总受影响行数 GET DIAGNOSTICS update_rows = ROW_COUNT; -- 统计实际插入的行数(通过查询目标表的新增数据,需根据业务主键判断) SELECT COUNT(*) INTO insert_rows FROM {{ this }} WHERE created_at >= CURRENT_TIMESTAMP() - INTERVAL '5 minutes'; -- 假设created_at是数据插入时的时间戳 -- 更新日志表 MERGE INTO log_table lt USING (SELECT '{{ table_name }}' AS TABLE_LOADED, insert_rows AS INSERTED_ROWS, (update_rows - insert_rows) AS UPDATED_ROWS) src ON lt.TABLE_LOADED = src.TABLE_LOADED WHEN MATCHED THEN UPDATE SET INSERTED_ROWS = src.INSERTED_ROWS, UPDATED_ROWS = src.UPDATED_ROWS, LOADED_AT = NOW() WHEN NOT MATCHED THEN INSERT (TABLE_LOADED, INSERTED_ROWS, UPDATED_ROWS) VALUES (src.TABLE_LOADED, src.INSERTED_ROWS, src.UPDATED_ROWS); END $$;
步骤2:模型配置
models: your_project: your_merge_model: materialized: incremental incremental_strategy: merge post-hook: - "{{ insert_log_sql }}" -- 直接引用上述SQL片段,或写入inline SQL
关键说明
- 不同数据库获取Merge行数的方式不同,需根据使用的数据库调整逻辑
- 如果需要复用日志逻辑,建议封装成存储过程或DBT宏,避免重复代码
- 需确保执行DBT的账号拥有操作日志表、调用存储过程(如果使用)的权限
内容的提问来源于stack exchange,提问作者Robertino Bonora
相关产品推荐
相关产品推荐

