如何在Airflow 2.9中每月将DAG调度至当月第一个工作日?
Airflow 每月第一个工作日调度DAG的实现方案
问题场景
需要将DAG调度到当月第一个工作日执行,比如2024年的调度时间示例:
2024-01-01 2024-02-01 2024-03-01 2024-04-01 2024-05-01 2024-06-03 -----> 这里是3号因为1号是周六
已了解Airflow的Timetable实现方式但觉得复杂,询问是否有其他可选方案,同时提供了自己的第一个工作日判断逻辑:
import pendulum year = 2024 for month in range(1,13): first_day_month = pendulum.DateTime(year=year,month=month,day=1) working_day = first_day_month if first_day_month.weekday() >= 0 and first_day_month.weekday() <= 4 else first_day_month.next(pendulum.MONDAY) print(working_day)
注:原代码中month从0开始会报错,已修正为从1到12
可选实现方式
自定义Timetable不是唯一的解决办法,以下两种方案更简便:
1. 基于@monthly触发+分支判断
核心思路是让DAG每月1号先触发,再通过分支任务判断当天是否为工作日:
- 如果是工作日,直接执行后续业务任务
- 如果是周末,计算出当月第一个工作日,触发DAG在该日期重新运行
代码示例:
from airflow import DAG from airflow.operators.python import BranchPythonOperator from airflow.operators.empty import EmptyOperator from airflow.operators.trigger_dagrun import TriggerDagRunOperator import pendulum from datetime import timedelta default_args = { 'owner': 'airflow', 'start_date': pendulum.datetime(2024, 1, 1), 'retries': 0 } def check_first_workday(**context): exec_date = context['execution_date'] # 判断是否为周一到周五 if exec_date.weekday() in range(0,5): return 'execute_business_tasks' else: # 计算第一个工作日:周六跳周一,周日跳周一 first_workday = exec_date.next(pendulum.MONDAY) if exec_date.weekday() ==5 else exec_date.add(days=1) context['ti'].xcom_push(key='target_date', value=first_workday) return 'trigger_dag_on_workday' with DAG( 'monthly_first_workday_dag', default_args=default_args, schedule_interval='@monthly', catchup=False ) as dag: check_workday = BranchPythonOperator( task_id='check_first_workday', python_callable=check_first_workday, provide_context=True ) # 替换成你的实际业务任务 execute_business_tasks = EmptyOperator(task_id='execute_business_tasks') trigger_dag = TriggerDagRunOperator( task_id='trigger_dag_on_workday', trigger_dag_id='monthly_first_workday_dag', execution_date="{{ ti.xcom_pull(key='target_date') }}", reset_dag_run=True ) end = EmptyOperator(task_id='end', trigger_rule='none_failed_min_one_success') check_workday >> [execute_business_tasks, trigger_dag] >> end
2. 简化版自定义Timetable
如果想要更原生的调度体验,不用每次触发分支判断,可以基于Airflow自带的CronExpressionTimetable扩展,只修改日期调整逻辑,比完全自定义Timetable代码量少很多:
from airflow.timetables.base import DagRunInfo, DataInterval from airflow.timetables.interval import CronExpressionTimetable from pendulum import DateTime class FirstWorkdayTimetable(CronExpressionTimetable): def __init__(self): # 先基于每月1号生成初始调度时间 super().__init__("0 0 1 * *", timezone="UTC") def next_dagrun_info( self, last_automated_dagrun: DateTime | None, restriction: DateTime | None, ) -> DagRunInfo | None: # 获取Cron生成的初始调度信息 base_info = super().next_dagrun_info(last_automated_dagrun, restriction) if not base_info: return None start_date = base_info.data_interval.start # 调整到第一个工作日 if start_date.weekday() >=5: adjusted_date = start_date.next(start_date.MONDAY) if start_date.weekday() ==5 else start_date.add(days=1) new_interval = DataInterval(start=adjusted_date, end=adjusted_date) return DagRunInfo(run_after=adjusted_date, data_interval=new_interval) return base_info
使用时直接在DAG中指定timetable=FirstWorkdayTimetable()即可。
总结
- 快速实现选方案1:无需深入Timetable底层,用现有Operator组合即可完成
- 追求原生调度体验选方案2:基于现有Timetable扩展,代码量少且符合Airflow调度逻辑
内容的提问来源于stack exchange,提问作者moth
相关产品推荐
相关产品推荐

