Airflow DAG中on_failure_callback无法正常触发的问题咨询
问题分析与解决方法
从你给出的报错信息和代码来看,核心问题出在on_failure_callback函数的参数定义和调用不匹配上。报错TypeError: on_failure_callback() takes 0 positional arguments but 1 was given说明Airflow在触发回调时传递了一个context参数,但你的回调函数实际没有接收这个参数的能力。
为什么会出现这个问题?
虽然你提供的代码里on_failure_callback定义了**context参数,但根据报错信息,实际运行的代码里这个函数大概率是没有参数的(比如误写成了def on_failure_callback():)。这种情况下,Airflow尝试传递context参数时就会触发类型错误。
正确的解决步骤
修正回调函数的参数定义
Airflow的on_failure_callback会强制传递一个包含任务上下文的参数,你可以用两种方式正确接收:- 直接接收
context对象:def on_failure_callback(context): # 可以通过context获取任务ID、执行时间等信息 print(f"Fail works! Task {context['ti'].task_id} failed at {context['execution_date']}") - 使用关键字参数(
**context)接收:def on_failure_callback(**context): print(f"Fail works! Task {context['ti'].task_id} failed")
- 直接接收
确保DAG默认参数正确引用函数
检查default_args里的on_failure_callback是指向函数对象,而不是字符串:default_args={ "owner": "airflow", "start_date": datetime(2020, 11, 1), "retries": 1, "retry_delay": timedelta(minutes=1), 'on_failure_callback': on_failure_callback, # 这里是函数对象,不是字符串 }验证测试场景
你在read_csv里故意写错的pd.DataFramee会正常触发任务失败,只要回调函数参数正确,就会打印出Fail works!的提示。
额外注意点
- 你的代码里
PythonOperator设置了provide_context=True,这是正确的,确保任务函数能接收context参数,但这和on_failure_callback的参数要求是独立的。 - 如果你需要在回调里做更多操作(比如发送告警邮件、记录失败日志),可以通过
context对象获取任务的详细信息,比如context['ti'].log_url可以拿到任务日志的链接。
内容的提问来源于stack exchange,提问作者Soumil Nitin Shah
相关产品推荐
相关产品推荐

