如何通过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
相关产品推荐
相关产品推荐

