Airflow 2.5:@task_group()无法接收params参数的解决方案咨询
Airflow 2.5中DAG params无法传递给task_group的解决方案
问题本质
@task_group()装饰的是任务组容器,并非可执行的任务实例,Airflow不会自动将DAG定义的params注入到task_group的定义函数中——只有@task装饰的可执行任务,才会被自动注入params参数,这就是直接在task_group里调用params会报错的原因。
解决方法
方法1:将params传递给task_group内部的任务
task_group内部定义的@task任务可以正常接收DAG的params,直接在内部任务中使用即可:
from datetime import datetime from airflow.decorators import dag, task, task_group @dag(dag_id='start-matrix', params={"a" : 1, "b" : 2}, schedule_interval=None, start_date=datetime(2021, 4, 5, 15, 0)) def startMatrix(): @task() def mytask(params=None): a = params["a"] print(f"mytask获取到a: {a}") @task_group() def mygroup(): # 在task_group内部定义任务,直接接收params @task() def inner_task(params=None): a = params["a"] b = params["b"] print(f"inner_task获取到a: {a}, b: {b}") inner_task() mytask() mygroup() startMatrix()
方法2:通过上下文获取DAG的params
如果需要在task_group的定义层面直接访问params,可以使用get_current_context()方法从Airflow上下文中提取:
from datetime import datetime from airflow.decorators import dag, task, task_group from airflow.operators.python import get_current_context @dag(dag_id='start-matrix', params={"a" : 1, "b" : 2}, schedule_interval=None, start_date=datetime(2021, 4, 5, 15, 0)) def startMatrix(): @task() def mytask(params=None): a = params["a"] print(f"mytask获取到a: {a}") @task_group() def mygroup(): # 从上下文获取params context = get_current_context() dag_params = context["params"] a = dag_params["a"] print(f"mygroup层面获取到a: {a}") # 内部任务依然可以正常接收params @task() def inner_task(params=None): b = params["b"] print(f"inner_task获取到b: {b}") inner_task() mytask() mygroup() startMatrix()
内容的提问来源于stack exchange,提问作者user1943569
相关产品推荐
相关产品推荐

