如何将DBT Core各模型执行日志写入SQL Server数据库表?
针对SQL Server的DBT模型执行日志记录方案
方案一:自定义宏+运行钩子(推荐)
利用DBT的钩子机制和自定义宏,直接在SQL Server中记录每个模型的执行数据:
- 初始化日志表
先在SQL Server中创建日志表,也可以用DBT模型来管理:
CREATE TABLE dbt_execution_logs ( log_id INT IDENTITY(1,1) PRIMARY KEY, project_name VARCHAR(100) NOT NULL, model_name VARCHAR(200) NOT NULL, start_time DATETIME2 NOT NULL, end_time DATETIME2, status VARCHAR(20) NOT NULL, -- 可选值:success/failed/skipped run_id VARCHAR(50) NOT NULL -- 关联同一次DBT运行的所有模型 );
- 编写日志记录宏
在DBT项目的macros目录下创建log_execution.sql:
{% macro log_model_execution(model, status, start_time, end_time) %} INSERT INTO dbt_execution_logs ( project_name, model_name, start_time, end_time, status, run_id ) VALUES ( '{{ project_name }}', '{{ model.unique_id }}', '{{ start_time }}', {% if end_time %}'{{ end_time }}'{% else %}NULL{% endif %}, '{{ status }}', '{{ invocation_id }}' ); {% endmacro %}
- 配置模型钩子
在dbt_project.yml中添加模型级钩子,捕获执行开始和成功完成的事件:
models: your_project_name: +pre-hook: "{{ log_model_execution(this, 'started', run_started_at, None) }}" +post-hook: "{{ log_model_execution(this, 'success', run_started_at, current_timestamp()) }}"
针对执行失败的模型,补充全局失败钩子(因为模型失败时不会触发post-hook):
on-run-failure: - "{{ log_failed_models(invocation_id) }}"
对应的log_failed_models宏可以读取target/run_results.json文件,提取失败模型信息后插入日志表。
方案二:解析DBT运行结果文件
DBT每次运行后会在target/目录生成run_results.json,可以编写脚本解析该文件并同步到SQL Server:
示例Python脚本片段:
import json import pyodbc from datetime import datetime # 读取DBT运行结果 with open('target/run_results.json', 'r') as f: run_data = json.load(f) # 连接SQL Server conn = pyodbc.connect('DRIVER={ODBC Driver 17 for SQL Server};SERVER=your_server;DATABASE=your_db;UID=user;PWD=password') cursor = conn.cursor() # 批量插入日志数据 for result in run_data['results']: project_name = run_data['metadata']['project_name'] model_name = result['unique_id'] start_time = datetime.fromisoformat(result['timing'][0]['started_at'].replace('Z', '+00:00')) end_time = datetime.fromisoformat(result['timing'][1]['completed_at'].replace('Z', '+00:00')) status = result['status'] cursor.execute(""" INSERT INTO dbt_execution_logs (project_name, model_name, start_time, end_time, status, run_id) VALUES (?, ?, ?, ?, ?, ?) """, (project_name, model_name, start_time, end_time, status, run_data['metadata']['invocation_id'])) conn.commit() cursor.close() conn.close()
可以将这个脚本加入到DBT运行的后续步骤(比如用shell脚本先执行dbt run,再运行该Python脚本)。
关于dbt_artifacts的补充
dbt_artifacts目前仅支持你提到的几种云数据平台,没有SQL Server适配。如果想基于它改造,需要调整宏中的SQL语法(比如身份列、日期函数),适配SQL Server的数据类型和系统表,不过这需要一定的DBT宏开发经验。
内容的提问来源于stack exchange,提问作者djohnjohn
相关产品推荐
相关产品推荐

