如何设置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
相关产品推荐
相关产品推荐

