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

Airflow:为工作日与周末设置不同的调度间隔

实现Airflow DAG工作日/周末差异化调度

方法1:使用Cron表达式组合

Airflow支持在schedule_interval中用逗号分隔多个Cron表达式,分别定义工作日和周末的调度规则:

  • 工作日(周一至周五)每小时运行:0 * * * 1-5
  • 周末(周六至周日)每3小时运行:0 */3 * * 6-0(Cron中周日可用0或7)

组合后的DAG配置示例:

from airflow import DAG
from datetime import datetime

default_args = {
    'start_date': datetime(2024, 1, 1),
}

with DAG(
    'data_warehouse_refresh',
    default_args=default_args,
    schedule_interval='0 * * * 1-5,0 */3 * * 6-0',
    catchup=False
) as dag:
    # 此处定义你的任务逻辑
    pass

注:多个Cron规则独立生效,Airflow会分别解析并触发任务,无需担心优先级问题。

方法2:自定义Timetable(Airflow 2.2+推荐)

如果需要更灵活的调度逻辑(比如后续扩展特殊日期规则),可以自定义Timetable类:

from airflow.timetables.base import DagRunInfo, DataInterval, Timetable
from airflow.utils.timezone import datetime as timezone_datetime
from croniter import croniter
import pendulum

class WeekdayWeekendTimetable(Timetable):
    def infer_manual_data_interval(self, run_after: timezone_datetime) -> DataInterval:
        # 手动触发时的时间区间逻辑,按需调整
        return DataInterval(start=run_after.subtract(hours=1), end=run_after)

    def next_dagrun_info(
        self,
        last_automated_data_interval: DataInterval | None,
        restriction: "TimeRestriction",
    ) -> DagRunInfo | None:
        if not last_automated_data_interval:
            # 首次运行起始时间
            start_time = restriction.earliest or pendulum.now().floor('hour')
        else:
            start_time = last_automated_data_interval.end

        while True:
            # 判断当前时间所属时段,选择对应Cron规则
            if start_time.weekday() in range(0,5):  # 周一至周五(0=周一,4=周五)
                cron_expr = '0 * * * 1-5'
                interval_hours = 1
            else:  # 周六(5)、周日(6)
                cron_expr = '0 */3 * * 6,0'
                interval_hours = 3
            
            next_run = croniter(cron_expr, start_time).get_next(timezone_datetime)
            if restriction.latest and next_run > restriction.latest:
                return None
            return DagRunInfo(
                run_after=next_run,
                data_interval=DataInterval(start=next_run.subtract(hours=interval_hours), end=next_run)
            )

# 应用自定义Timetable
with DAG(
    'data_warehouse_refresh',
    default_args=default_args,
    timetable=WeekdayWeekendTimetable(),
    catchup=False
) as dag:
    # 任务定义
    pass

验证调度规则

可以通过Airflow命令行验证下一次执行时间:

airflow dags next-execution data_warehouse_refresh

也可在Airflow UI的DAG详情页查看「Next Run」字段,确认规则是否生效。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 08:15:46