You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 10:17:32