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

Airflow分支任务:特定条件下依赖上游任务的实现问题

解决Airflow任务依赖与条件执行的问题

核心思路是结合ShortCircuitOperator控制任务启停,同时通过触发规则处理不同场景下的依赖逻辑,完美匹配你的需求:

实现代码

from airflow import DAG
from airflow.operators.python import PythonOperator, ShortCircuitOperator
from airflow.utils.trigger_rule import TriggerRule
from datetime import datetime

def run_task_a_func(**context):
    # 替换为task_a的实际业务逻辑,可通过context获取专属输入参数
    task_a_params = context["dag_run"].conf.get("task_a_params", {})
    print(f"执行task_a,参数:{task_a_params}")

def run_task_b_func(**context):
    # 替换为task_b的清理逻辑
    print("执行task_b,清理失效产物")

def should_run_task_a(**context):
    # 根据DAG参数判断是否执行task_a
    return context["dag_run"].conf.get("run_task_a", False)

def should_run_task_b(**context):
    # 根据DAG参数判断是否执行task_b
    return context["dag_run"].conf.get("run_task_b", False)

with DAG(
    dag_id="task_dependency_dag",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    # 控制task_a是否执行的前置检查
    check_task_a = ShortCircuitOperator(
        task_id="check_task_a",
        python_callable=should_run_task_a,
        provide_context=True
    )

    task_a = PythonOperator(
        task_id="task_a",
        python_callable=run_task_a_func,
        provide_context=True
    )

    # 控制task_b是否执行的前置检查,关键设置触发规则
    check_task_b = ShortCircuitOperator(
        task_id="check_task_b",
        python_callable=should_run_task_b,
        provide_context=True,
        trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS
    )

    task_b = PythonOperator(
        task_id="task_b",
        python_callable=run_task_b_func,
        provide_context=True
    )

    # 构建依赖关系
    check_task_a >> task_a >> check_task_b
    check_task_b >> task_b

关键逻辑说明

  1. ShortCircuitOperator的作用:

    • check_task_a根据run_task_a参数决定是否执行后续的task_a,参数为False时直接跳过task_a。
    • check_task_b根据run_task_b参数决定是否启动task_b,参数为False时直接跳过task_b。
  2. 触发规则的关键设置:
    给check_task_b设置TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS,意味着:

    • 当task_a执行成功时,check_task_b会等待task_a完成后再执行。
    • 当task_a被跳过(run_task_a=False)时,只要上游没有失败任务,check_task_b依然会正常执行,确保task_b可以独立启动。

测试场景验证

  • 场景1:run_task_a=True且run_task_b=True → task_a执行完成后,task_b自动启动。
  • 场景2:run_task_a=False且run_task_b=True → task_a被跳过,task_b直接启动。
  • 场景3:run_task_b=False → 无论run_task_a状态如何,task_b都不会执行。

内容的提问来源于stack exchange,提问作者dwolfeu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 09:53:18