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

如何用TaskFlow API创建含每日/每月条件任务的Airflow DAG?

使用Airflow TaskFlow API实现每日+月度任务调度

方案一:序列结构(非1号跳过task_2)

利用@task.short_circuit装饰器,仅在每月1号时允许task_2执行,非1号时直接跳过task_2,task_3不受影响继续执行。这种结构保持线性依赖,逻辑更直观。

from datetime import datetime
from airflow.decorators import dag, task
from airflow.operators.empty import EmptyOperator

@dag(
    start_date=datetime(2024, 9, 9),
    schedule_interval="@daily",  # 每日调度
    catchup=False
)
def daily_monthly_seq_dag():
    task_1 = EmptyOperator(task_id='task_1')

    @task.short_circuit(task_id="check_monthly_run")
    def should_run_task_2(**context):
        # 获取执行日期的天
        execution_day = context['execution_date'].day
        return execution_day == 1

    task_2 = EmptyOperator(task_id='task_2')
    task_3 = EmptyOperator(task_id='task_3')

    # 依赖配置:非1号时短路任务跳过task_2,直接触发task_3
    task_1 >> should_run_task_2() >> task_2 >> task_3
    should_run_task_2() >> task_3

daily_monthly_seq_dag()

方案二:分支结构(始终执行task_1和task_3,仅1号执行task_2)

通过@task.branch实现分支逻辑,同时给task_3设置触发规则,保证无论分支走向task_2还是直接指向task_3,task_3都会正常执行。

from datetime import datetime
from airflow.decorators import dag, task
from airflow.operators.empty import EmptyOperator
from airflow.utils.trigger_rule import TriggerRule

@dag(
    start_date=datetime(2024, 9, 9),
    schedule_interval="@daily",
    catchup=False
)
def daily_monthly_branch_dag():
    task_1 = EmptyOperator(task_id='task_1')

    @task.branch(task_id="monthly_branch")
    def branch_logic(**context):
        execution_day = context['execution_date'].day
        if execution_day == 1:
            return ['task_2', 'task_3']  # 1号同时触发task_2和task_3
        else:
            return ['task_3']  # 非1号直接触发task_3

    task_2 = EmptyOperator(task_id='task_2')
    # 设置触发规则:只要上游至少一个任务成功就执行
    task_3 = EmptyOperator(task_id='task_3', trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS)

    # 依赖配置
    task_1 >> branch_logic()
    branch_logic() >> task_2 >> task_3
    branch_logic() >> task_3

daily_monthly_branch_dag()

关键说明

  1. 方案一短路逻辑:@task.short_circuit返回True时执行后续的task_2;返回False时task_2被标记为跳过,通过补充的依赖直接触发task_3,确保每日执行。
  2. 方案二分支与触发规则:分支任务可返回多个任务ID实现并行触发;task_3的触发规则保证无论上游是分支直接触发,还是task_2执行完成,只要没有失败且至少一个任务成功,就会运行。
  3. 两个方案都配置schedule_interval="@daily"实现每日调度,catchup=False避免回溯执行历史任务。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 17:35:15