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

如何配置Airflow DAG每月最后一个周二执行调度?

实现Airflow DAG每月最后一个周二运行的配置方法

方法一:Cron表达式+任务内日期校验

先用基础Cron表达式匹配所有周二,再在起始任务中校验当前日期是否为当月最后一个周二,不符合则跳过后续任务。

  1. 配置DAG的schedule_interval为每周二运行:
schedule_interval="0 0 * * 2"  # 每周二0点触发
  1. 用@task.short_circuit实现日期校验逻辑:
from airflow.decorators import dag, task
import pendulum

@dag(
    start_date=pendulum.datetime(2024, 1, 1, tz="UTC"),
    schedule_interval="0 0 * * 2",
    catchup=False
)
def last_tuesday_dag():
    @task.short_circuit
    def validate_last_tuesday():
        today = pendulum.now("UTC").date()
        # 获取当月最后一天
        month_last_day = today.end_of("month").date()
        # 计算当月最后一个周二:周二对应weekday=1,倒推天数
        days_to_subtract = (month_last_day.weekday() - 1) % 7
        month_last_tuesday = month_last_day - pendulum.duration(days=days_to_subtract)
        return today == month_last_tuesday

    @task
    def core_task():
        print("执行每月最后一个周二的核心任务")

    validate_last_tuesday() >> core_task()

last_tuesday_dag()

short_circuit任务返回False时会跳过后续任务,仅当当天是当月最后一个周二时才执行核心逻辑。

方法二:自定义调度器(推荐)

通过自定义Timetable类让Airflow直接计算出准确的运行时间,无需额外校验,调度逻辑更精准。

from airflow.decorators import dag, task
import pendulum
from airflow.timetables.base import DagRunInfo, DataInterval, TimeRestriction
from airflow.timetables.interval import CronDataIntervalTimetable
from typing import Optional

class LastTuesdayTimetable(CronDataIntervalTimetable):
    def next_dagrun_info(
        self,
        last_automated_dagrun: Optional[DagRun],
        restriction: TimeRestriction,
    ) -> Optional[DagRunInfo]:
        # 确定计算起始时间
        if last_automated_dagrun:
            start_calc = last_automated_dagrun.execution_date.add(months=1)
        else:
            start_calc = restriction.earliest or pendulum.now("UTC")
        
        # 计算当月最后一个周二
        month_last_day = start_calc.end_of("month")
        days_to_subtract = (month_last_day.weekday() - 1) % 7
        target_date = month_last_day.subtract(days=days_to_subtract).replace(hour=0, minute=0, second=0)
        
        # 校验是否在时间限制范围内
        if restriction.latest and target_date > restriction.latest:
            return None
        
        return DagRunInfo(
            execution_date=target_date,
            data_interval=DataInterval(start=target_date.start_of("month"), end=target_date)
        )

@dag(
    start_date=pendulum.datetime(2024, 1, 1, tz="UTC"),
    timetable=LastTuesdayTimetable("0 0 * * 2"),
    catchup=False
)
def last_tuesday_dag():
    @task
    def core_task():
        print("执行每月最后一个周二的核心任务")

    core_task()

last_tuesday_dag()

这种方式让调度器直接生成符合要求的运行计划,是Airflow官方推荐的灵活调度实现方式。

方法三:分支任务实现跳过逻辑

用BranchPythonOperator实现分支判断,选择执行核心任务或跳过任务:

from airflow import DAG
from airflow.operators.python import BranchPythonOperator, PythonOperator
from airflow.utils.dates import days_ago
import pendulum

def check_last_tuesday(**context):
    exec_date = context["execution_date"].date()
    month_last_day = exec_date.end_of("month").date()
    days_to_subtract = (month_last_day.weekday() - 1) % 7
    month_last_tuesday = month_last_day - pendulum.duration(days=days_to_subtract)
    return "core_task" if exec_date == month_last_tuesday else "skip_task"

def core_task(**context):
    print("执行每月最后一个周二的核心任务")

def skip_task(**context):
    print("今日非当月最后一个周二,跳过任务")

with DAG(
    dag_id="last_tuesday_dag",
    start_date=days_ago(1),
    schedule_interval="0 0 * * 2",
    catchup=False
) as dag:
    branch_check = BranchPythonOperator(
        task_id="check_last_tuesday",
        python_callable=check_last_tuesday,
        provide_context=True
    )

    core = PythonOperator(
        task_id="core_task",
        python_callable=core_task
    )

    skip = PythonOperator(
        task_id="skip_task",
        python_callable=skip_task
    )

    branch_check >> [core, skip]

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 14:05:24