如何在Airflow中调整任务执行顺序与分支任务数量
解决方案
核心思路是先定义task_2_a、task_2_b、task_2_c的固定顺序依赖,再通过BranchPythonOperator返回对应分支的起始任务,最后让所有task_2系列任务的终点都指向task_3,配合合适的触发规则实现需求。
步骤1:修改分支逻辑函数
将原本返回多个任务的逻辑改为返回分支的起始任务,任务间的顺序依赖会自动触发后续任务:
def get_path(**kwargs): params = kwargs.get('params', {}) path = params.get('path') if path == '1': return 'task_2_a' elif path == '2': return 'task_2_b' elif path == '3': return 'task_2_c' elif path == '4': return 'task_2_a' # 执行a后自动触发b else: return 'task_2_a' # 执行a→b→c依次触发
步骤2:设置正确的任务依赖
先定义task_2系列的顺序依赖,再连接主流程和分支,最后配置task_3的触发规则以适配不同分支的执行情况:
# 定义task_2系列的顺序执行关系 task_2_a >> task_2_b >> task_2_c # 主流程依赖 task_1 >> branch_1 # 分支指向各起始任务 branch_1 >> [task_2_a, task_2_b, task_2_c] # 配置task_3的触发规则:只要至少一个前置任务成功且无失败就执行 task_3 = SQLExecuteQueryOperator( task_id='task_3', sql=f""" insert into {{{{dag_run.conf.name}}}} (some_text) select some_text from (select '333' as some_text) as tab """, trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS ) # 连接task_3到complete [task_2_a, task_2_b, task_2_c] >> task_3 >> complete
完整修改后的DAG代码
from datetime import datetime, timedelta, date from airflow import DAG from airflow.operators.python import BranchPythonOperator from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator from airflow.models import DagRun from airflow.utils.trigger_rule import TriggerRule from airflow.operators.dummy import DummyOperator def get_path(**kwargs): params = kwargs.get('params', {}) path = params.get('path') if path == '1': return 'task_2_a' elif path == '2': return 'task_2_b' elif path == '3': return 'task_2_c' elif path == '4': return 'task_2_a' else: return 'task_2_a' with DAG( 'test', description='test', tags=["test"], schedule_interval=None, start_date=datetime(2025, 7, 1), default_args={ 'retries': 0, 'retry_delay': timedelta(minutes=1), 'conn_id': 'sgk_gp' }, params={ 'name':'', 'path':'' } ) as dag: task_1 = SQLExecuteQueryOperator( task_id='task_1', sql=f""" drop table if exists {{{{dag_run.conf.name}}}}; create table {{{{dag_run.conf.name}}}} ( some_text character varying ) """ ) branch_1 = BranchPythonOperator( task_id='branch_1', python_callable=get_path, provide_context=True, do_xcom_push=False ) task_2_a = SQLExecuteQueryOperator( task_id='task_2_a', sql=f""" insert into {{{{dag_run.conf.name}}}} (some_text) select some_text from (select 'aaa' as some_text) as tab """ ) task_2_b = SQLExecuteQueryOperator( task_id='task_2_b', sql=f""" insert into {{{{dag_run.conf.name}}}} (some_text) select some_text from (select 'bbb' as some_text) as tab """ ) task_2_c = SQLExecuteQueryOperator( task_id='task_2_c', sql=f""" insert into {{{{dag_run.conf.name}}}} (some_text) select some_text from (select 'ccc' as some_text) as tab """ ) task_3 = SQLExecuteQueryOperator( task_id='task_3', sql=f""" insert into {{{{dag_run.conf.name}}}} (some_text) select some_text from (select '333' as some_text) as tab """, trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS ) complete = DummyOperator(task_id="complete", trigger_rule=TriggerRule.NONE_FAILED) # 配置所有依赖关系 task_2_a >> task_2_b >> task_2_c task_1 >> branch_1 >> [task_2_a, task_2_b, task_2_c] [task_2_a, task_2_b, task_2_c] >> task_3 >> complete
逻辑验证
- 分支1(path=1):
task_1→branch_1→task_2_a→task_3→complete - 分支2(path=2):
task_1→branch_1→task_2_b→task_3→complete - 分支3(path=3):
task_1→branch_1→task_2_c→task_3→complete - 分支4(path=4):
task_1→branch_1→task_2_a→task_2_b→task_3→complete(严格顺序执行) - 分支5(path=其他):
task_1→branch_1→task_2_a→task_2_b→task_2_c→task_3→complete(严格顺序执行)
内容的提问来源于stack exchange,提问作者Serge Kashlik
相关产品推荐
相关产品推荐

