如何配置Airflow DAG每月最后一个周二执行调度?
实现Airflow DAG每月最后一个周二运行的配置方法
方法一:Cron表达式+任务内日期校验
先用基础Cron表达式匹配所有周二,再在起始任务中校验当前日期是否为当月最后一个周二,不符合则跳过后续任务。
- 配置DAG的
schedule_interval为每周二运行:
schedule_interval="0 0 * * 2" # 每周二0点触发
- 用
@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
相关产品推荐
相关产品推荐

