为何将on_failure_callback放入dag_default_args中无法生效?
Airflow中
on_failure_callback放入独立dag_default_args失效原因 两段代码逻辑几乎一致,唯一区别是第一段将on_failure_callback直接写在DAG的default_args字面量字典中,回调能正常触发;第二段把包含on_failure_callback的default_args抽成模块级的dag_default_args变量,回调就失效了。
生效代码示例
#datetime from datetime import timedelta, datetime from airflow.operators.python import PythonVirtualenvOperator import pendulum from airflow import DAG import logging def on_failure_callback(context): logging.error("enter on_failure_callback-----------") # python callable function def print_hello(): raise Exception("helloworld exception") with DAG( dag_id="hello_world_dag", default_args={ 'on_failure_callback': on_failure_callback, }, schedule_interval="00 02 * * *", start_date=pendulum.datetime(2022, 12, 1, tz="UTC"), catchup=False, max_active_runs=1, ) as dag: export = PythonVirtualenvOperator( python_callable=print_hello, task_id="hello", dag=dag, )
失效代码示例
#datetime from datetime import timedelta, datetime from airflow.operators.python import PythonVirtualenvOperator import pendulum from airflow import DAG import logging def on_failure_callback(context): logging.error("enter on_failure_callback-----------") dag_default_args = { 'retries': 1, 'on_failure_callback': on_failure_callback, } # python callable function def print_hello(): raise Exception("helloworld exception") with DAG( dag_id="hello_world_dag", default_args=dag_default_args, schedule_interval="00 02 * * *", start_date=pendulum.datetime(2022, 12, 1, tz="UTC"), catchup=False, max_active_runs=1, ) as dag: export = PythonVirtualenvOperator( python_callable=print_hello, task_id="hello", dag=dag, )
失效原因分析
核心问题出在Airflow的DAG序列化机制:
- Airflow需要将DAG配置序列化后存储到元数据库,供调度器、执行器等组件读取使用。
- Python自定义函数属于不可直接序列化的对象。当
on_failure_callback被放在模块级的dag_default_args字典中时,这个字典会在DAG文件加载阶段被Airflow的序列化逻辑处理,函数对象无法被正确持久化,最终导致任务失败时找不到有效的回调函数。 - 而直接在DAG的
default_args参数中定义字面量字典时,Airflow会在DAG解析的上下文里保留函数的有效引用,序列化逻辑能正确识别并关联这个回调函数,因此可以正常触发。
内容的提问来源于stack exchange,提问作者seaguest
相关产品推荐
相关产品推荐

