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

如何在Airflow中按日期跳过特定任务?

单个DAG中按日期跳过指定任务的可行方案

针对你要在单个DAG中混合每日、每周任务的需求,以下是三种实用方案,无需修改任务列表或设置任务级schedule_interval:

方案一:用ShortCircuitOperator控制任务执行

通过ShortCircuitOperator为每个任务(组)添加前置判断,符合条件则继续执行后续任务,不符合则直接跳过。这种方式逻辑清晰,任务独立性强。

from airflow import DAG
from airflow.operators.python import ShortCircuitOperator, PythonOperator
from datetime import datetime, timedelta

def daily_task_1():
    print("执行每日任务1")

def daily_task_2():
    print("执行每日任务2")

def weekly_task_1():
    print("执行每周任务1(周一运行)")

def weekly_task_2():
    print("执行每周任务2(周五运行)")

def check_daily(**context):
    # 每日任务无条件放行
    return True

def check_weekly_monday(**context):
    # 判断执行日期是否为周一(isoweekday():1=周一,7=周日)
    return context['execution_date'].isoweekday() == 1

def check_weekly_friday(**context):
    return context['execution_date'].isoweekday() == 5

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

with DAG(
    'mixed_schedule_dag',
    default_args=default_args,
    schedule_interval='@daily',  # DAG按每日调度
    catchup=False,
) as dag:
    # 每日任务分支
    check_daily_op = ShortCircuitOperator(
        task_id='check_daily',
        python_callable=check_daily,
        provide_context=True
    )
    daily_1 = PythonOperator(task_id='daily_task_1', python_callable=daily_task_1)
    daily_2 = PythonOperator(task_id='daily_task_2', python_callable=daily_task_2)

    # 每周一任务分支
    check_monday_op = ShortCircuitOperator(
        task_id='check_monday',
        python_callable=check_weekly_monday,
        provide_context=True
    )
    weekly_1 = PythonOperator(task_id='weekly_task_1', python_callable=weekly_task_1)

    # 每周五任务分支
    check_friday_op = ShortCircuitOperator(
        task_id='check_friday',
        python_callable=check_weekly_friday,
        provide_context=True
    )
    weekly_2 = PythonOperator(task_id='weekly_task_2', python_callable=weekly_task_2)

    # 依赖设置:三个分支并行执行
    check_daily_op >> [daily_1, daily_2]
    check_monday_op >> weekly_1
    check_friday_op >> weekly_2

方案二:在任务内部嵌入日期判断逻辑

直接在任务的业务代码开头添加日期判断,不符合条件则提前退出,不执行核心逻辑。这种方式无需额外控制任务,代码更紧凑。

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta

def daily_task_1():
    print("执行每日任务1")

def daily_task_2():
    print("执行每日任务2")

def weekly_task_1(**context):
    exec_date = context['execution_date']
    if exec_date.isoweekday() != 1:
        print("今日非周一,跳过每周任务1")
        return
    # 核心业务逻辑
    print("执行每周任务1(周一运行)")

def weekly_task_2(**context):
    exec_date = context['execution_date']
    if exec_date.isoweekday() != 5:
        print("今日非周五,跳过每周任务2")
        return
    # 核心业务逻辑
    print("执行每周任务2(周五运行)")

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

with DAG(
    'mixed_schedule_dag_v2',
    default_args=default_args,
    schedule_interval='@daily',
    catchup=False,
) as dag:
    daily_1 = PythonOperator(task_id='daily_task_1', python_callable=daily_task_1)
    daily_2 = PythonOperator(task_id='daily_task_2', python_callable=daily_task_2)
    weekly_1 = PythonOperator(task_id='weekly_task_1', python_callable=weekly_task_1, provide_context=True)
    weekly_2 = PythonOperator(task_id='weekly_task_2', python_callable=weekly_task_2, provide_context=True)

    # 所有任务并行执行,各自判断是否运行
    [daily_1, daily_2, weekly_1, weekly_2]

方案三:用BranchPythonOperator集中控制分支

通过BranchPythonOperator在一个任务中集中判断所有需要执行的任务,返回符合条件的task_id列表,Airflow会自动执行这些任务,跳过未被选中的任务。适合需要统一管理执行条件的场景。

from airflow import DAG
from airflow.operators.python import BranchPythonOperator, PythonOperator
from datetime import datetime, timedelta

def daily_task_1():
    print("执行每日任务1")

def daily_task_2():
    print("执行每日任务2")

def weekly_task_1():
    print("执行每周任务1(周一运行)")

def weekly_task_2():
    print("执行每周任务2(周五运行)")

def decide_tasks(**context):
    exec_date = context['execution_date']
    tasks_to_run = ['daily_task_1', 'daily_task_2']  # 每日任务默认执行
    if exec_date.isoweekday() == 1:
        tasks_to_run.append('weekly_task_1')
    if exec_date.isoweekday() == 5:
        tasks_to_run.append('weekly_task_2')
    return tasks_to_run

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

with DAG(
    'mixed_schedule_dag_v3',
    default_args=default_args,
    schedule_interval='@daily',
    catchup=False,
) as dag:
    branch_op = BranchPythonOperator(
        task_id='decide_which_tasks_to_run',
        python_callable=decide_tasks,
        provide_context=True
    )
    daily_1 = PythonOperator(task_id='daily_task_1', python_callable=daily_task_1)
    daily_2 = PythonOperator(task_id='daily_task_2', python_callable=daily_task_2)
    weekly_1 = PythonOperator(task_id='weekly_task_1', python_callable=weekly_task_1)
    weekly_2 = PythonOperator(task_id='weekly_task_2', python_callable=weekly_task_2)

    # 依赖设置:分支任务指向所有可能的任务
    branch_op >> [daily_1, daily_2, weekly_1, weekly_2]

方案对比

  • ShortCircuitOperator:逻辑拆分清晰,每个任务的执行条件独立,适合任务间无强依赖的场景。
  • 任务内部判断:代码简洁,无需额外控制任务,适合简单的条件判断场景。
  • BranchPythonOperator:集中管理分支逻辑,便于统一修改执行条件,但任务较多时需注意维护返回的task_id列表,避免遗漏或错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 18:43:12