Airflow动态任务映射:基于非直接前置上游任务输出创建映射任务的问题
Airflow 动态任务映射的非直接上游依赖问题
Airflow 2.7官方文档明确动态任务映射基于前置任务的输出创建,但未覆盖如何基于任意上游任务(非直接前置)生成映射任务的场景。
比如要实现执行链 A >> B >> C[],其中映射任务C[]的数据源是A的输出,但按常规写法会自动生成多余的A >> C[]直接依赖(图形视图中可见该不必要的连线)。
问题示例代码
from airflow.decorators import dag, task, task_group from datetime import datetime import random as rnd @task def A(): return [x for x in range(int(rnd.random()*10)+1)] @task def B(): print("intermediate") @task def C(x): return x + 10 @dag(dag_id="TEST_DYNAMIC_SIMPLE",start_date=datetime(2023, 10, 11), catchup=False,schedule=None,tags=["TESTS"]) def TEST_DYNAMIC_SIMPLE(): tA = A() exC = C.expand(x=tA) tA >> B() >> exC globals()["TEST_DYNAMIC_SIMPLE"] = TEST_DYNAMIC_SIMPLE() # 本地测试 if __name__ == "__main__": dag = globals()['TEST_DYNAMIC_SIMPLE'] dag.test()
解决方案
Airflow的expand方法会自动将输入参数对应的任务设为映射任务的上游,这就是A >> C[]直接依赖产生的原因。要实现预期的依赖链,可通过以下方式处理:
方法1:新增中间传递任务
创建一个仅传递A输出的中间任务,让该任务依赖B,再让C[]基于中间任务的输出展开。这样C[]的上游仅为中间任务,而中间任务同时依赖A和B,既保证执行顺序,又避免了A到C的直接连线。
修复后的代码:
from airflow.decorators import dag, task, task_group from datetime import datetime import random as rnd @task def A(): return [x for x in range(int(rnd.random()*10)+1)] @task def B(): print("intermediate") @task def pass_a_output(a_data): # 仅传递A的输出,无额外业务逻辑 return a_data @task def C(x): return x + 10 @dag(dag_id="TEST_DYNAMIC_FIXED",start_date=datetime(2023, 10, 11), catchup=False,schedule=None,tags=["TESTS"]) def TEST_DYNAMIC_FIXED(): tA = A() tB = B() # 中间任务绑定A的输出,同时设置B为其上游 pass_task = pass_a_output(a_data=tA) tB >> pass_task # 映射任务基于中间任务的输出展开,上游仅为pass_task exC = C.expand(x=pass_task) globals()["TEST_DYNAMIC_FIXED"] = TEST_DYNAMIC_FIXED() # 本地测试 if __name__ == "__main__": dag = globals()['TEST_DYNAMIC_FIXED'] dag.test()
说明
虽然C[]必须等待A完成才能获取输出(A是隐式依赖),但通过中间任务可以让DAG的图形视图更贴合预期的A >> B >> C[]依赖链,同时保证执行逻辑正确。
内容的提问来源于stack exchange,提问作者Aleksei Meledin
相关产品推荐
相关产品推荐

