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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 22:27:41