如何触发dbt在测试WARN/Fail时执行模型及connection报错解决
解决dbt Macro的"connection未定义"错误与外部SQL执行需求
一、修复"local variable 'connection' referenced before assignment"错误
错误原因很直接:你在macro里直接调用了connection变量,但没有提前初始化它。在dbt中,正确获取数据库连接的方式有两种:
1. 用dbt内置adapter初始化连接
在macro开头添加连接定义:
{% macro check_test_results() %} {% set connection = adapter.get_connection() %} {# 你的后续逻辑 #} {% endmacro %}
2. 用dbt_utils的get_connection宏(需先安装dbt_utils)
如果项目依赖dbt_utils包,也可以这么写:
{% macro check_test_results() %} {% set connection = dbt_utils.get_connection() %} {# 你的后续逻辑 #} {% endmacro %}
注意:数据库操作完成后,一定要加{% do adapter.close_connection(connection) %}关闭连接,避免连接池耗尽。
二、替换硬编码SQL为执行外部SQL文件
dbt提供了两种方式调用外部SQL文件,根据你的场景选择:
方式1:用{% include %}直接引入外部SQL
假设你的告警逻辑写在macros/test_alert_logic.sql里,直接在macro中引入执行:
{% macro check_test_results() %} {% set connection = adapter.get_connection() %} {% if execute %} -- 查询测试状态 {% set test_status_query %} select status from dbt_test_results where ... {# 你的状态过滤逻辑 #} {% endset %} {% set results = run_query(test_status_query) %} -- 遍历结果,触发告警 {% for row in results %} {% if row.status in ['warn', 'fail'] %} {% set alert_sql = include('macros/test_alert_logic.sql') %} {% do run_query(alert_sql) %} {% endif %} {% endfor %} {% endif %} {% do adapter.close_connection(connection) %} {% endmacro %}
方式2:把外部SQL封装成独立macro调用
如果告警逻辑复杂,把它写成单独的macro(比如macros/run_alert_model.sql):
{% macro run_alert_model() %} -- 这里是外部SQL的具体逻辑 insert into alert_logs(test_name, status, alert_time) select test_name, status, current_timestamp from dbt_test_results where status in ('warn', 'fail'); {% endmacro %}
然后在主macro里直接调用:
{% macro check_test_results() %} {% set connection = adapter.get_connection() %} {% if execute %} {% set test_status_query %} select test_name, status from dbt_test_results where ... {% endset %} {% set results = run_query(test_status_query) %} {% set has_alert = results.columns.status.values() | select('in', ['warn', 'fail']) | list | length > 0 %} {% if has_alert %} {% do run_query(run_alert_model()) %} {% endif %} {% endif %} {% do adapter.close_connection(connection) %} {% endmacro %}
三、完整可运行示例
结合以上两点,完整的macro代码:
{% macro check_test_results() %} {% set connection = adapter.get_connection() %} {% if execute %} -- 查询最近1天的测试结果 {% set test_status_query %} select status from dbt_test_results where execution_time >= current_date - interval '1 day' {% endset %} {% set results = run_query(test_status_query) %} -- 判断是否需要触发告警 {% set has_alert = results.columns.status.values() | select('in', ['warn', 'fail']) | list | length > 0 %} {% if has_alert %} -- 执行外部SQL文件中的告警逻辑 {% set alert_logic = include('macros/test_failure_alert.sql') %} {% do run_query(alert_logic) %} {{ log("已触发告警模型执行", info=True) }} {% endif %} {% endif %} {% do adapter.close_connection(connection) %} {% endmacro %}
关键注意事项
- 外部SQL文件路径要正确,dbt从项目根目录开始查找
- 所有数据库操作必须放在
{% if execute %}块里,避免解析阶段执行无效操作 - 务必关闭连接,防止连接资源泄漏
内容的提问来源于stack exchange,提问作者JuniorAngelo
相关产品推荐
相关产品推荐

