求Different non-linear dag execution非线性架构的相关实现示例
非线性DAG执行架构示例参考
架构匹配说明
你所需的是多根节点并行+单分支内线性执行的非对称DAG架构,4个根节点dag1/dag2/dag3/dag4的depend_on=true属性代表根节点仅校验自身历史执行状态、互相之间无前置依赖,对应拓扑结构如下:
dag1->dag11 | dag2->dag21->dag22 | dag3->dag31 | dag4
可运行代码示例(Airflow 2.x)
以下是符合要求的DAG定义代码:
from airflow import DAG from airflow.operators.dummy import DummyOperator from datetime import datetime with DAG( dag_id="non_linear_multi_branch_sample", start_date=datetime(2024, 1, 1), schedule_interval=None, catchup=False ) as dag: # 分支1 线性执行 dag1 = DummyOperator(task_id="dag1", depend_on_past=True) dag11 = DummyOperator(task_id="dag11") dag1 >> dag11 # 分支2 线性执行 dag2 = DummyOperator(task_id="dag2", depend_on_past=True) dag21 = DummyOperator(task_id="dag21") dag22 = DummyOperator(task_id="dag22") dag2 >> dag21 >> dag22 # 分支3 线性执行 dag3 = DummyOperator(task_id="dag3", depend_on_past=True) dag31 = DummyOperator(task_id="dag31") dag3 >> dag31 # 分支4 无后续节点 dag4 = DummyOperator(task_id="dag4", depend_on_past=True)
执行规则说明
- 调度触发后4个根任务同时并行启动,每个根任务仅会在自身上一次调度执行成功的情况下才会运行本次实例
- 各分支内部严格按线性顺序执行,前序任务执行成功才会触发后续任务
- 不同分支的执行进度互不干扰,各自满足依赖即可推进
- 若使用其他DAG调度框架(如DolphinScheduler、Prefect),仅需按上述拓扑结构配置任务依赖,开启根节点的自依赖属性即可实现完全相同的执行逻辑
内容的提问来源于stack exchange,提问作者LeComp
相关产品推荐
相关产品推荐

