如何在Airflow中移除上下游任务依赖?
Airflow 任务依赖移除方案
背景:Airflow 任务依赖定义方式
在Airflow的DAG中,定义两个任务的依赖关系有三种常见方式:
首先定义基础任务:
from airflow.operators.dummy import DummyOperator t1 = DummyOperator(task_id='dummy_1') t2 = DummyOperator(task_id='dummy_2')
依赖定义的三种写法:
# 方式A:使用 >> 操作符 t1 >> t2 # 方式B:调用set_upstream方法 t2.set_upstream(t1) # 方式C:调用set_downstream方法 t1.set_downstream(t2)
问题描述
是否支持在依赖定义完成后,移除任务的上游或下游依赖?
在动态生成大规模任务及依赖的DAG场景中,任务创建完成后需要调整依赖关系——比如插入新任务到现有任务链中,同时移除原有任务间的依赖。
举例说明:
已有任务t1和t2,且t1 >> t2,现在希望添加新任务t3,将依赖关系改为t1 >> t3 >> t2,同时移除t1与t2的原有依赖,是否可行?
对应示例代码:
from airflow import DAG from airflow.operators.dummy import DummyOperator from datetime import datetime def function_that_creates_dags_dynamically(): tasks = { 't1': DummyOperator(task_id='dummy_1'), 't2': DummyOperator(task_id='dummy_2'), } tasks['t1'] >> tasks['t2'] return tasks with DAG( dag_id='test_dag', start_date=datetime(2021, 1, 1), catchup=False, tags=['example'], ) as dag: tasks = function_that_creates_dags_dynamically() t3 = DummyOperator(task_id='dummy_3') tasks['t1'] >> t3 t3 >> tasks['t2'] # 如何移除tasks['t1']与tasks['t2']的依赖?
解决方案
Airflow的任务对象(继承自BaseOperator)维护了用于管理依赖的集合属性,可直接操作这些属性移除已有依赖:
具体操作方式
要移除t1和t2之间的原有依赖,只需执行以下两行代码:
# 移除t1的下游依赖t2 tasks['t1'].downstream_task_ids.remove(tasks['t2'].task_id) # 移除t2的上游依赖t1 tasks['t2'].upstream_task_ids.remove(tasks['t1'].task_id)
也可以直接操作任务实例列表:
# 移除t1下游任务列表中的t2 tasks['t1'].downstream_list.remove(tasks['t2']) # 移除t2上游任务列表中的t1 tasks['t2'].upstream_list.remove(tasks['t1'])
完整修改后的代码
from airflow import DAG from airflow.operators.dummy import DummyOperator from datetime import datetime def function_that_creates_dags_dynamically(): tasks = { 't1': DummyOperator(task_id='dummy_1'), 't2': DummyOperator(task_id='dummy_2'), } tasks['t1'] >> tasks['t2'] return tasks with DAG( dag_id='test_dag', start_date=datetime(2021, 1, 1), catchup=False, tags=['example'], ) as dag: tasks = function_that_creates_dags_dynamically() t3 = DummyOperator(task_id='dummy_3') tasks['t1'] >> t3 t3 >> tasks['t2'] # 移除t1与t2的原有依赖 if tasks['t2'].task_id in tasks['t1'].downstream_task_ids: tasks['t1'].downstream_task_ids.remove(tasks['t2'].task_id) if tasks['t1'].task_id in tasks['t2'].upstream_task_ids: tasks['t2'].upstream_task_ids.remove(tasks['t1'].task_id)
注意事项
- 操作集合时,若目标task_id不存在会抛出
KeyError,建议先通过in判断存在性再执行移除操作; - 该方案适用于Airflow 2.x及以上版本,修改后DAG的依赖关系会被正确解析,不影响任务调度。
内容的提问来源于stack exchange,提问作者Tokyo
相关产品推荐
相关产品推荐

