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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 05:11:12