Airflow中能否同时配置DAG级与任务级on_failure_callback?
Airflow DAG级与任务级on_failure_callback的执行逻辑
核心结论
任务级的on_failure_callback会覆盖DAG级的回调,不会同时触发。具体执行逻辑分两种场景:
- 当单个任务失败时:如果任务自身配置了
on_failure_callback,只会执行该任务的回调;若任务未配置,则依次 fallback 到default_args中的回调、DAG级的回调。 - 当DAG本身出现调度层面的失败(无任务执行就宣告失败):只会触发DAG级的
on_failure_callback。
优先级规则
Airflow回调的优先级从高到低为:
- 任务级单独配置的
on_failure_callback default_args中定义的on_failure_callback- DAG对象直接配置的
on_failure_callback
结合你的示例说明
你的DAG定义中,default_args包含独立的失败回调,同时DAG本身配置了failure_callback_dag:
with DAG( default_args={"on_failure_callback": 自定义任务默认回调}, dag_id = "test_dag", schedule_interval=None, start_date=datetime.datetime(2021, 1, 1), catchup=False, on_failure_callback=failure_callback_dag ) as dag: # 任务定义
- 如果任务未单独配置回调,会优先执行
default_args里的回调,而非DAG级的failure_callback_dag - 如果任务自己指定了
on_failure_callback,则只会执行该任务的回调
实现两者都触发的方法
如果需要任务失败时同时执行任务级和DAG级的回调,不能依赖Airflow默认逻辑,需手动在任务级回调中调用DAG级回调:
def task_failure_callback(context): # 执行任务级失败操作 handle_task_failure(context) # 调用DAG级失败回调 failure_callback_dag(context) # 任务配置示例 test_task = PythonOperator( task_id="test_task", python_callable=your_task_func, on_failure_callback=task_failure_callback )
内容的提问来源于stack exchange,提问作者Deshana Desai
相关产品推荐
相关产品推荐

