如何基于节假日日历实现Airflow DAG差异化时段调度
解决方案
推荐方案(Airflow 2.2及以上版本):自定义Timetable + 节假日数据预同步
该方案是目前业界处理非标准调度场景的最优实现,可完全避开你提到的所有问题:
- 第一步:预同步节假日数据到本地
编写一个独立的轻量DAG,每天凌晨执行一次,从业务数据库拉取最新的节假日列表,存储到所有scheduler、worker节点都能访问的共享存储路径(如NAS、共享磁盘),格式为JSON即可。完全避免业务DAG运行或解析时频繁查询业务数据库,彻底解决DBA侧的压力问题。
同步逻辑示例:import json import os from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime HOLIDAY_FILE_PATH = "/shared/holidays.json" def sync_holidays(): # 这里替换为你查询数据库获取节假日的逻辑 holidays = query_holiday_from_db() with open(HOLIDAY_FILE_PATH, 'w', encoding='utf-8') as f: json.dump([d.strftime("%Y-%m-%d") for d in holidays], f) with DAG( dag_id='sync_holiday_calendar', schedule_interval='0 1 * * *', catchup=False, default_args={'start_date': datetime(2024,1,1)} ) as dag: sync_task = PythonOperator( task_id='sync_holidays', python_callable=sync_holidays ) - 第二步:实现自定义Timetable
Airflow 2.2推出的Timetable能力就是为了解决标准cron无法覆盖的调度场景,你可以自定义调度逻辑,仅在需要的时候生成DAG Run,不会产生多余的运行记录。
自定义Timetable示例:from airflow.timetables.base import Timetable, DataInterval, DagRunInfo from datetime import datetime, time, timedelta import json HOLIDAY_FILE_PATH = "/shared/holidays.json" class HolidayAdjustableTimetable(Timetable): def infer_manual_data_interval(self, run_after: datetime) -> DataInterval: start = run_after.replace(hour=0, minute=0, second=0, microsecond=0) end = start.replace(hour=23, minute=59, second=59) return DataInterval(start=start, end=end) def next_dagrun_info( self, *, last_automated_data_interval: DataInterval | None, restriction, ) -> DagRunInfo | None: # 计算下一个待调度的日期 if last_automated_data_interval is None: next_date = datetime.now(tz=restriction.timezone).date() else: next_date = last_automated_data_interval.start.date() + timedelta(days=1) # 跳过周末 while next_date.weekday() >=5: next_date += timedelta(days=1) # 读取本地节假日文件判断 with open(HOLIDAY_FILE_PATH, 'r', encoding='utf-8') as f: holidays = set(json.load(f)) date_str = next_date.strftime("%Y-%m-%d") run_time = time(9,0) if date_str in holidays else time(17,0) run_at = datetime.combine(next_date, run_time, tzinfo=restriction.timezone) return DagRunInfo.interval(start=run_at, end=run_at) - 第三步:业务DAG引用自定义Timetable
with DAG( dag_id='jobname', default_args=args, timetable=HolidayAdjustableTimetable(), catchup=False, ) as dag: # 直接编写正常的任务逻辑,无需额外添加时间/节假日判断 your_task = PythonOperator( task_id='run_your_business', python_callable=your_task_func )
该方案的优势:
- 仅保留1个业务DAG,无冗余DAG
- 仅生成符合条件的DAG Run,无多余的失败/跳过记录,GUI展示干净
- 所有节假日判定读本地文件,零业务数据库查询压力
兼容方案(Airflow <2.2版本):优化版双调度+ShortCircuitOperator
如果你的Airflow版本较低不支持Timetable,可以在你第三个方案的基础上做两点优化即可:
- 节假日判定读提前预同步的本地JSON文件,不再直接查业务数据库
- 用
ShortCircuitOperator替代直接抛出异常,不满足运行条件时标记为skipped状态,不会产生失败记录,也可配置为不在GUI中展示跳过的运行记录
代码示例:
from airflow.operators.python import ShortCircuitOperator from datetime import datetime import json HOLIDAY_FILE_PATH = "/shared/holidays.json" def check_run_condition(): with open(HOLIDAY_FILE_PATH, 'r', encoding='utf-8') as f: holidays = set(json.load(f)) today = datetime.now().date() current_hour = datetime.now().hour is_holiday = today.strftime("%Y-%m-%d") in holidays return (is_holiday and current_hour ==9) or (not is_holiday and current_hour ==17) with DAG( dag_id='jobname', default_args=args, schedule_interval='0 9,17 * * 1-5', catchup=False, ) as dag: check_condition = ShortCircuitOperator( task_id='check_run_condition', python_callable=check_run_condition, ignore_downstream_trigger_rules=False ) your_task = PythonOperator( task_id='run_your_business', python_callable=your_task_func ) check_condition >> your_task
内容的提问来源于stack exchange,提问作者user3240688
相关产品推荐
相关产品推荐

