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

如何设置Airflow数据区间以上月首日为起始、上月末日为结束?

如何设置Airflow数据区间以上月首日为起始、上月末日为结束?

我有一个ETL流程,每月7日上午8点执行,需要处理上月首日至上月末日的数据:

  • 3月7日处理2月1日至2月28日(或29日)的数据
  • 4月7日处理3月1日至3月31日的数据

当前使用CronDataIntervalTimetable(cron="0 8 7 * *")时,数据区间不符合预期(比如3月7日执行时得到2023-02-07 8:00 - 2023-03-07 8:00),以下是两种可行的解决方法:


方案一:用Jinja模板宏直接计算目标日期

无需修改调度规则,直接通过Airflow内置的时间宏计算上月的起始和结束日期,代码改动最小:

from __future__ import annotations
from datetime import datetime
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.timetables.interval import CronDataIntervalTimetable

def do_the_etl(start_date_str: str, end_date_str: str) -> None:
    start_date = datetime.fromisoformat(start_date_str)
    end_date = datetime.fromisoformat(end_date_str)
    print(f"[Fake ETL] Querying for the period {start_date}-{end_date}")

# 利用execution_date计算上月首日和末日
ETL_DATE_START = "{{ execution_date.start_of('month').subtract(months=1).isoformat() }}"
# 上月末日 = 本月首日减1秒,确保覆盖上月最后一秒的数据
ETL_DATE_END = "{{ execution_date.start_of('month').subtract(seconds=1).isoformat() }}"

with DAG(
    dag_id="etl_test",
    start_date=datetime(2023, 2, 4),
    schedule_interval=CronDataIntervalTimetable(cron="0 8 7 * *", timezone="Etc/UTC"),
    catchup=False,
) as dag:
    run_etl = PythonOperator(
        task_id="etl",
        python_callable=do_the_etl,
        op_kwargs={
            "start_date_str": ETL_DATE_START,
            "end_date_str": ETL_DATE_END,
        },
    )

run_etl

逻辑说明:

  • execution_date是当前DAG的执行时间(即每月7日8点)
  • start_of('month')获取当月1日0点,减1个月得到上月1日0点
  • 当月1日0点减1秒,即为上月最后一天的23:59:59,完美覆盖上月全量数据

方案二:自定义Timetable(推荐,使数据区间语义化)

如果希望DAG的data_interval_start和data_interval_end本身就对应上月完整区间,可以自定义Timetable,让调度时间与数据区间解耦:

1. 实现自定义Timetable

from datetime import datetime
from airflow.timetables.base import DagRunInfo, DataInterval, TimeRestriction
from airflow.timetables.interval import CronTimetable
from pendulum import DateTime

class LastMonthTimetable(CronTimetable):
    def __init__(self, cron: str, timezone: str):
        super().__init__(cron, timezone)

    def infer_data_interval(self, run_after: DateTime) -> DataInterval:
        # 计算当月首日
        current_month_start = run_after.start_of('month')
        # 上月首日作为区间起始,当月首日作为区间结束(左闭右开)
        return DataInterval(
            start=current_month_start.subtract(months=1),
            end=current_month_start
        )

    def next_dagrun_info(
        self,
        last_automated_dagrun: DateTime | None,
        restriction: TimeRestriction,
    ) -> DagRunInfo | None:
        dag_run_info = super().next_dagrun_info(last_automated_dagrun, restriction)
        if not dag_run_info:
            return None
        # 绑定调度时间与对应的数据区间
        data_interval = self.infer_data_interval(dag_run_info.run_after)
        return DagRunInfo(run_after=dag_run_info.run_after, data_interval=data_interval)

2. 在DAG中使用自定义Timetable

from __future__ import annotations
from datetime import datetime
from airflow import DAG
from airflow.operators.python import PythonOperator

def do_the_etl(start_date_str: str, end_date_str: str) -> None:
    start_date = datetime.fromisoformat(start_date_str)
    end_date = datetime.fromisoformat(end_date_str)
    print(f"[Fake ETL] Querying for the period {start_date}-{end_date}")

# 直接使用语义化的data_interval变量
ETL_DATE_START = "{{ data_interval_start.isoformat() }}"
ETL_DATE_END = "{{ data_interval_end.subtract(seconds=1).isoformat() }}"

with DAG(
    dag_id="etl_test_custom_timetable",
    start_date=datetime(2023, 2, 4),
    schedule_interval=LastMonthTimetable(cron="0 8 7 * *", timezone="Etc/UTC"),
    catchup=False,
) as dag:
    run_etl = PythonOperator(
        task_id="etl",
        python_callable=do_the_etl,
        op_kwargs={
            "start_date_str": ETL_DATE_START,
            "end_date_str": ETL_DATE_END,
        },
    )

run_etl

逻辑说明:

  • 自定义Timetable后,每月7日8点的调度任务,对应的data_interval_start为上月1日0点,data_interval_end为本月1日0点
  • 因为Airflow数据区间是左闭右开规则,若ETL查询用< end_date,可直接使用data_interval_end;若需要包含上月最后一秒,减去1秒即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 21:45:03