Airflow Task Group执行顺序异常:如何按指定流程配置DAG?
问题原因及解决方案
为什么task_3a会立即执行
你的DAG存在两个核心问题导致task_3a提前执行:
- 依赖链断裂:分支任务
branching_task仅关联了task_2,未将task_3纳入分支逻辑,且task_3没有明确的上游依赖(原代码中task_3仅被task_2触发,当分支不选择task_2时,task_3无依赖会直接启动),进而导致后续的task_group也提前执行。 - task_group任务串行:你在
task_after_3_before_4_group中返回task_3a() >> task_3b(),这会让两个任务串行执行,不符合你需要并行的要求。
调整后的完整代码
以下是符合你指定执行顺序的修改版代码:
import logging from airflow.decorators import task, dag, task_group from airflow.utils.dates import days_ago @dag( dag_id='taskflow_conditional_dag', start_date=days_ago(1), schedule_interval=None, catchup=False, ) def my_dag(): logger = logging.getLogger("airflow.task") @task def task_1(): logger.info("Task 1 running") return "run_task_" @task.branch def branching_task(data): # 条件满足时执行task_2,否则跳过task_2 if data == "run_task_2": return "task_2" # 返回空列表表示不执行任何分支任务(即跳过task_2) return [] @task def task_2(): logger.info("Task 2 running") pass @task( trigger_rule="none_failed_or_skipped" # 不管task_2是成功还是被跳过,都执行task_3 ) def task_3(): logger.info("Task 3 running") pass @task def task_4(): logger.info("Task 4 running") pass @task_group() def task_after_3_before_4_group(): @task def task_3a(): logger.info("Task 3a running.") pass @task def task_3b(): logger.info("Task 3b running.") pass # 返回任务列表表示并行执行 return [task_3a(), task_3b()] @task def task_5(): logger.info("Task 5 running") pass # 任务实例化 data = task_1() decision = branching_task(data) task_2_result = task_2() task_3_result = task_3() task_group_parallel = task_after_3_before_4_group() task_4_result = task_4() task_5_result = task_5() # 构建依赖链 data >> decision decision >> task_2_result # task_3同时依赖task_1和task_2,trigger_rule确保跳过task_2时仍能执行 [data, task_2_result] >> task_3_result # task_3完成后启动并行的3a和3b task_3_result >> task_group_parallel # 最后依次执行task_4和task_5 task_group_parallel >> task_4_result >> task_5_result dag = my_dag()
关键修改说明
- 修复分支逻辑:
branching_task在不需要执行task_2时返回空列表,标记task_2为跳过状态。 - 调整task_3的触发规则:设置
trigger_rule="none_failed_or_skipped",确保无论task_2是成功执行还是被跳过,task_3都会启动。 - 修正依赖关系:让
task_3同时依赖task_1和task_2,保证执行顺序符合要求。 - 实现任务并行:在task_group中返回任务列表
[task_3a(), task_3b()],让两个任务并行执行。 - 简化trigger_rule:
task_4和task_5无需额外触发规则,默认的all_success即可满足(因为并行任务全部完成后才会触发下游)。
内容的提问来源于stack exchange,提问作者Scott
相关产品推荐
相关产品推荐

