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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 00:55:13