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

如何通过Airflow的DbtRunOperationOperator获取dbt宏的返回值?

问题描述

定义了如下dbt宏,预期返回Hello, dbt!:

{% macro return_hello() %}
    {% set result = 'Hello, dbt!' %}
    {{ return(result) }}
{% endmacro %}

使用Airflow的DbtRunOperationOperator执行该宏,设置do_xcom_push=True,并通过以下代码尝试获取返回值:

def fetch_result(context):
     ti = context['ti']
     value = ti.xcom_pull(task_ids='dbt_task')
     logging.info(value)

任务定义如下:

print_task = DbtRunOperationOperator(task_id='dbt_task', macro='return_hello', do_xcom_push=True)

实际运行后,日志中打印的是dbt任务的元数据JSON,而非预期的Hello, dbt!。

解决方案

1. 修改dbt宏,将结果输出到标准输出

宏的return值不会被Airflow主动捕获,需要将结果主动输出到stdout:

{% macro return_hello() %}
    {% set result = 'Hello, dbt!' %}
    {{ result }}  {# 将结果输出到标准输出,供Airflow捕获 #}
    {{ return(result) }}
{% endmacro %}

2. 调整Airflow任务参数,开启输出捕获

在DbtRunOperationOperator中添加fetch_output=True参数,该参数会让Operator捕获宏执行时输出到stdout的内容,并将其作为XCom的值推送,而非默认的任务元数据:

print_task = DbtRunOperationOperator(
    task_id='dbt_task',
    macro='return_hello',
    do_xcom_push=True,
    fetch_output=True  # 开启宏输出捕获
)

原理说明

DbtRunOperationOperator默认推送到XCom的是dbt任务运行的元数据(如执行状态、耗时等)。通过fetch_output=True,Operator会捕获宏执行过程中输出到标准输出的内容,替换默认的元数据作为XCom返回值。配合宏中主动输出结果的代码,即可获取到预期的Hello, dbt!。

内容的提问来源于stack exchange,提问作者Aaryan Ohekar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 17:52:31