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
相关产品推荐
相关产品推荐

