Apache Airflow DAG触发异常:触发上一周期而非当前周期
问题解答
这是Apache Airflow的预期行为,核心源于Airflow的调度模型设计:任务会在一个调度周期结束后触发,此时任务对应的ds(即execution_date)是该周期的起始时间,而非触发时刻的当前时间:
- 分钟级调度(
*/1 * * * *):15:08触发的任务,对应15:07-15:08的周期,因此ds为15:07; - 年度调度(
0 12 1 1 *):2023年1月1日触发的任务,对应2022年1月1日-2023年1月1日的周期,因此ds为2022-01-01。
这种设计是为了适配绝大多数数据处理场景——比如天级任务在凌晨执行,处理前一天的全量数据,ds可直接对应待处理的日期。
适配方案(获取当前周期时间)
如果你的业务需要使用触发时刻对应的当前周期时间,可通过以下两种方式调整:
1. 直接使用上下文变量next_execution_date
Airflow会在provide_context=True的前提下,自动将next_execution_date传入任务函数,这个变量代表当前任务周期的结束时间,也就是你期望的“当前周期”时间。
以你提供的第一个DAG为例,修改test_print函数:
def test_print(ds, next_execution_date, foo, **kwargs): # 按需求格式化next_execution_date为字符串 now = next_execution_date.strftime('%Y-%m-%d %H:%M:%S') data2send = {'the_date_n_hour': now} # 其余逻辑保持不变
2. 基于ds手动计算当前周期时间
如果无法直接使用next_execution_date,可根据调度间隔,从ds推导当前周期时间:
- 分钟级调度:
datetime.strptime(ds, '%Y-%m-%d %H:%M:%S') + timedelta(minutes=1) - 年度调度:
datetime.strptime(ds, '%Y-%m-%d') + relativedelta(years=1)
以年度DAG的get_holidays函数为例:
def get_holidays(ds, gtp_id, **kwargs): # 计算当前年度的起始和结束时间 current_year_start = (datetime.strptime(ds, '%Y-%m-%d') + relativedelta(years=1)).date() holi_start_date = str(current_year_start) holi_end_date = str(current_year_start + relativedelta(years=1)) # 其余逻辑保持不变
注意事项
- 若开启
catchup=True,历史任务会遵循同样的时间逻辑执行,需确保调整后的计算规则对历史周期也适用; - 不要混淆
execution_date(任务对应的周期起始)与start_date(DAG的初始运行日期)的概念。
内容的提问来源于stack exchange,提问作者Andrew A. O.
相关产品推荐
相关产品推荐

