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

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.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 09:40:29