Airflow:无需额外Operator读取UI触发DAG时传入的CLI输入失败问题
解决Airflow中从UI传递参数到DAG函数的问题
你的代码问题出在全局变量kpi='{{ kpi}}'的处理上——这个字符串在DAG解析阶段就被固定了,Airflow不会自动渲染它,所以你打印的只是模板字符串本身,而不是触发时传入的参数值。而且触发DAG时传入的参数其实是存在DAG Run的配置里的,我们可以直接从上下文的dag_run对象中获取,不需要依赖模板变量或者额外的Operator。
修复方案
核心修改点
- 移除全局的
kpi变量定义,因为它无法被正确渲染。 - 在
get_data_from_bq函数中,通过kwargs['dag_run'].conf直接读取触发时传入的参数。 - 确保
PythonOperator的provide_context=True(这会把Airflow的上下文对象传递给你的函数,包括dag_run)。
修改后的完整代码
from airflow import DAG from airflow.utils.dates import days_ago from airflow.operators.python_operator import PythonOperator from airflow import models from airflow.models import Variable from google.cloud import bigquery from airflow.configuration import conf LOCATION = Variable.get("HDM_PROJECT_LOCATION") PROJECT_ID = Variable.get("HDM_PROJECT_ID") client = bigquery.Client() # default arguments default_dag_args = { 'start_date': days_ago(0), 'retries': 0, 'project_id': PROJECT_ID } def get_data_from_bq(**kwargs): # 从DAG Run的配置中提取传入的kpi参数 # 如果没有传入kpi,默认返回空字符串或自定义默认值 dag_config = kwargs.get('dag_run', {}).conf or {} kpi_value = dag_config.get('kpi', '') print("op is:") print(kpi_value) with models.DAG( '00_test_sql1', schedule_interval=None, default_args=default_dag_args) as dag: v_run_sql_01 = PythonOperator( task_id='Run_SQL', provide_context=True, # 必须开启,才能获取dag_run等上下文对象 python_callable=get_data_from_bq, location=LOCATION, use_legacy_sql=False )
为什么原来的代码不生效?
- 全局变量的时机问题:DAG文件会被Airflow定期解析,全局变量
kpi='{{ kpi}}'在解析阶段就被赋值为字符串"{{ kpi}}",此时还没有触发DAG,所以模板不会被渲染。 - 参数的存储位置:从UI触发DAG时传入的参数,会被存储在当前DAG Run的
conf属性中,而不是直接作为模板变量存在全局上下文中。
这样修改后,当你触发DAG并传入{"kpi":"ID123"}时,函数里就能正确打印出ID123了。
内容的提问来源于stack exchange,提问作者Yug
相关产品推荐
相关产品推荐

