Airflow多链式Expand调用下动态任务映射的分组执行问题
解决方案:用动态生成TaskGroup实现分组串行执行
你说得对,直接用expand确实会导致所有task_2实例执行完毕后才会启动task_3,而要实现每个task_2完成后立刻触发对应的task_3,我们可以通过动态生成TaskGroup的方式来实现——虽然TaskGroup没有expand方法,但我们可以遍历task_1的返回结果,为每个task_num单独创建一个包含task_2和task_3的TaskGroup,让组内任务串行,组间任务并行。
下面是完整的实现代码:
from airflow import DAG from airflow.decorators import task, task_group from pendulum import datetime, now @task def task_1(): # 这里返回动态数量的任务标识,示例返回[0,1,2,3,4] return list(range(5)) @task def task_2(task_num): print(f"Executing task_2 for task_num: {task_num}") return task_num @task def task_3(task_num): print(f"Executing task_3 for task_num: {task_num}") return task_num @task_group def create_task_group(task_num): # 每个TaskGroup内,task_2完成后立即执行task_3 t2 = task_2(task_num=task_num) t3 = task_3(task_num=task_num) t2 >> t3 with DAG(dag_id="my_dag", start_date=now(), schedule_interval=None) as dag: # 先执行task_1获取动态任务数量列表 task_1_result = task_1() # 遍历task_1的返回结果,为每个task_num创建独立的TaskGroup for num in task_1_result: create_task_group(task_num=num)
代码解释:
- TaskGroup工厂函数:
create_task_group是一个被@task_group装饰的函数,接收task_num作为参数,内部定义了该组内的task_2和task_3,并通过>>设置了它们的上下游依赖,确保task_2完成后立刻启动task_3。 - 动态生成分组:在DAG定义中,我们先执行
task_1拿到返回的任务列表,然后遍历这个列表,为每个task_num实例化一个create_task_group——这样每个分组都是独立的,组内串行、组间并行,完全匹配你想要的任务图结构。 - 执行逻辑对比:和你之前的方案不同,这个实现不会等待所有
task_2完成,只要某个分组内的task_2执行完毕,对应的task_3就会立即启动,完美解决了你的需求。
内容的提问来源于stack exchange,提问作者bruno
相关产品推荐
相关产品推荐

