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

Apache Airflow 2.5.0:如何在Python VirtualEnv Operator中使用配置JSON

在Airflow 2.5.0的VirtualEnvOperator中获取DAG触发配置值

你之前尝试方法失败的原因

  • 方法1:未在VirtualEnvOperator中设置provide_context=True,且函数未接收**kwargs,导致上下文未传入,自然取不到dag_run
  • 方法2:虚拟环境内的代码不会被Airflow的Jinja引擎渲染,直接在函数里写模板语法无效
  • 方法3:设置DAG的params后,需通过kwargs['params']['conf1']访问,且同样需要开启provide_context=True才能拿到上下文
  • 方法4:手动创建的Jinja Template对象无法获取Airflow注入的上下文变量,所以识别不了dag_run

两种可行解决方案

方案一:通过上下文直接获取dag_run配置

from airflow import DAG
from airflow.operators.python import VirtualEnvOperator
from datetime import datetime

def process_trigger_conf(**kwargs):
    # 从上下文提取dag_run对象
    dag_run = kwargs.get('dag_run')
    if dag_run and dag_run.conf:
        # 获取conf1的值
        conf1_val = dag_run.conf.get('conf1')
        print(f"获取到conf1的值: {conf1_val}")
        # 这里可将值存入变量供后续逻辑使用
        return conf1_val
    print("未获取到触发配置")
    return None

with DAG(
    dag_id='fetch_trigger_conf',
    start_date=datetime(2023, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    fetch_conf_task = VirtualEnvOperator(
        task_id='fetch_conf',
        python_callable=process_trigger_conf,
        provide_context=True,  # 必须开启,才能传递上下文
        requirements=[],  # 根据你的依赖需求添加,比如pandas等
        system_site_packages=False
    )

方案二:通过op_kwargs传递渲染后的配置值

这种方式更直接,无需在函数里处理dag_run,直接把渲染好的配置值传入函数:

from airflow import DAG
from airflow.operators.python import VirtualEnvOperator
from datetime import datetime

def process_conf(conf1_val, **kwargs):
    print(f"直接拿到conf1的值: {conf1_val}")
    # 存入变量使用
    return conf1_val

with DAG(
    dag_id='fetch_conf_via_op_kwargs',
    start_date=datetime(2023, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    fetch_conf_task = VirtualEnvOperator(
        task_id='fetch_conf',
        python_callable=process_conf,
        op_kwargs={
            # 用Jinja模板渲染配置,Airflow会自动替换为触发时传入的值
            'conf1_val': "{{ dag_run.conf.get('conf1', '默认值') }}"
        },
        requirements=[],
        system_site_packages=False
    )

注意事项

  • 触发DAG时务必正确传入配置JSON:{"conf1": "test"}
  • 方案一中的provide_context=True是上下文传递的核心,不能省略
  • 方案二中的Jinja模板只能写在Operator的参数里,不能写在函数内部
  • 使用.get('conf1', '默认值')可以避免因未传入conf1导致的KeyError

内容的提问来源于stack exchange,提问作者s misra

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 21:00:58