如何在Airflow TaskFlow任务中获取dag_run上下文与配置信息
问题:TaskFlow任务中如何读取DAG启动配置?
DAG启动时携带的配置JSON如下:
{"foo" : "bar"}
原本使用PythonOperator读取该配置的代码:
my_task = PythonOperator( task_id="my_task", op_kwargs={"foo": "{{ dag_run.conf['foo'] }}"}, python_callable=lambda foo: print(foo))
现在要替换为TaskFlow任务,代码框架如下,需要实现读取配置值的逻辑:
@task def my_task: # 如何获取foo?
请问如何在此TaskFlow任务中获取context、dag_run,或是读取上述配置JSON中的值?
解决方案
有三种常用方法可以在TaskFlow任务中读取DAG启动配置:
方法1:直接通过模板参数注入
和PythonOperator的思路一致,直接把配置值作为参数传入TaskFlow任务,使用Jinja2模板语法:
from airflow.decorators import task @task def my_task(foo): print(foo) # 在DAG中调用任务时传入参数 my_task(foo="{{ dag_run.conf['foo'] }}")
方法2:获取完整上下文对象
如果需要访问更多上下文信息,可以在任务函数中添加**context参数,再从context中取出dag_run:
from airflow.decorators import task @task def my_task(**context): dag_run = context["dag_run"] foo = dag_run.conf.get("foo") print(foo)
方法3:直接声明获取DagRun对象
Airflow支持直接在TaskFlow任务中通过参数声明获取dag_run对象:
from airflow.decorators import task from airflow.models import DagRun @task def my_task(dag_run: DagRun): foo = dag_run.conf.get("foo") print(foo)
注意:如果配置中可能不存在foo字段,建议使用get()方法避免KeyError,也可以添加默认值,比如dag_run.conf.get("foo", "default_value")。
内容的提问来源于stack exchange,提问作者Aaron Brager
相关产品推荐
相关产品推荐

