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

Airflow:无需额外Operator读取UI触发DAG时传入的CLI输入失败问题

解决Airflow中从UI传递参数到DAG函数的问题

你的代码问题出在全局变量kpi='{{ kpi}}'的处理上——这个字符串在DAG解析阶段就被固定了,Airflow不会自动渲染它,所以你打印的只是模板字符串本身,而不是触发时传入的参数值。而且触发DAG时传入的参数其实是存在DAG Run的配置里的,我们可以直接从上下文的dag_run对象中获取,不需要依赖模板变量或者额外的Operator。

修复方案

核心修改点

  1. 移除全局的kpi变量定义,因为它无法被正确渲染。
  2. 在get_data_from_bq函数中,通过kwargs['dag_run'].conf直接读取触发时传入的参数。
  3. 确保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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 06:52:49