如何在Airflow中用单个DAG实现同逻辑的日/月双调度?
单个Airflow文件实现每日+每月调度的DAG复用方案
完全可以用单个文件实现,核心思路是在同一个Python文件中定义两个独立的DAG实例(对应每日、每月调度),复用同一套业务逻辑函数,只针对调度规则、日期参数生成逻辑做区分,避免重复维护代码。
完整实现代码(命名为pull_data_combined.py)
from airflow import DAG from airflow.operators.python import PythonOperator from pendulum import datetime, pendulum # 通用默认参数,两个DAG共用 default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2024, 1, 1), 'retries': 1 } # ---------------------- 复用的核心业务逻辑 ---------------------- def start_process(date): # 这里写入你的数据提取业务逻辑,比如连接服务器拉取对应日期的数据 print(f"开始处理日期 {date} 的数据") # 你的业务代码... # ---------------------- 每日调度DAG ---------------------- dag_daily = DAG( 'scheduler_daily', default_args=default_args, catchup=False, schedule_interval='0 0 * * *', description='每日调度:提取当日数据' ) # 每日任务:使用Airflow内置宏{{ ds }}获取执行日期(格式YYYY-MM-DD) pull_daily_task = PythonOperator( task_id="start_pulling_data_from_server", python_callable=start_process, op_kwargs={'date': '{{ ds }}'}, dag=dag_daily ) # ---------------------- 每月调度DAG ---------------------- dag_monthly = DAG( 'scheduler_monthly', default_args=default_args, catchup=False, schedule_interval='0 0 1 * *', description='每月调度:提取上月全月数据' ) # 生成上月的所有日期列表(DAG解析时自动计算) def generate_last_month_dates(): current_month_start = pendulum.now().start_of('month') last_month_start = current_month_start.subtract(months=1) last_month_end = current_month_start.subtract(days=1) date_list = [] current_date = last_month_start while current_date <= last_month_end: date_list.append(current_date.strftime('%Y-%m-%d')) current_date = current_date.add(days=1) return date_list # 动态生成上月每日的任务(每个任务对应一天的数据提取) for day in generate_last_month_dates(): pull_monthly_task = PythonOperator( task_id=f"start_pulling_data_from_server_{day}", python_callable=start_process, op_kwargs={'date': day}, dag=dag_monthly )
关键注意事项
- 复用业务逻辑:
start_process函数只写一次,两个DAG的任务都调用它,后续修改业务逻辑只需改这一处 - 每日任务用Airflow宏:用
{{ ds }}替代硬编码的date_now,更符合Airflow调度规范,自动匹配执行日期 - 每月任务唯一task_id:循环生成任务时,必须给每个任务设置唯一的
task_id(这里用日期后缀),否则Airflow会抛出任务ID重复的错误 - DAG自动识别:Airflow会扫描文件中的所有
DAG实例,自动加载scheduler_daily和scheduler_monthly两个调度任务,和原来两个文件的效果完全一致
内容的提问来源于stack exchange,提问作者andikapr
相关产品推荐
相关产品推荐

