如何使用TaskFlow API为on_failure指定DAG或任务失败回调函数
TaskFlow API配置DAG/任务失败回调的实现方案
TaskFlow API的回调配置和传统写法核心逻辑一致,仅在装饰器声明、DAG实例化阶段传入对应参数即可,以下是完整实现步骤:
1. 定义失败回调函数
回调函数必须接收context作为入参,这是Airflow传入的运行上下文对象,包含DAG ID、任务ID、执行时间等核心信息,示例代码如下:
def dag_failure_callback(context): # 自定义DAG失败逻辑,比如发送告警、打印错误信息等 dag_id = context.get("dag_id") execution_date = context.get("execution_date") print(f"DAG {dag_id} 运行失败,执行时间:{execution_date}") def task_failure_callback(context): # 自定义单个任务失败的回调逻辑 task_id = context.get("task_instance").task_id print(f"任务 {task_id} 运行失败")
2. 配置DAG全局失败回调
如果需要整个DAG运行最终状态为失败时触发回调,直接在@dag装饰器中传入on_failure_callback参数即可:
from airflow.decorators import dag, task from datetime import datetime @dag( schedule_interval=None, start_date=datetime(2024, 1, 1), # 配置DAG级失败回调 on_failure_callback=dag_failure_callback, catchup=False ) def my_taskflow_dag(): @task def test_task(): # 模拟任务失败 raise Exception("任务运行出错") test_task() dag = my_taskflow_dag()
3. 配置单个Task的失败回调
如果仅需要特定任务失败时触发回调,在@task装饰器中传入on_failure_callback参数即可:
@dag( schedule_interval=None, start_date=datetime(2024, 1, 1), catchup=False ) def my_taskflow_dag(): # 给该任务单独配置失败回调 @task(on_failure_callback=task_failure_callback) def risky_task(): raise Exception("风险任务出错") # 该任务出错不会触发上面的task_failure_callback @task def normal_task(): print("普通任务运行") risky_task() >> normal_task() dag = my_taskflow_dag()
4. 混合使用传统Operator的回调配置
如果TaskFlow DAG中混用了传统Operator(比如BashOperator、PythonOperator),直接在Operator实例化时传入on_failure_callback参数即可:
from airflow.operators.bash import BashOperator @dag( schedule_interval=None, start_date=datetime(2024, 1, 1), catchup=False ) def my_mixed_dag(): # 传统Operator配置回调 bash_task = BashOperator( task_id="bash_task", bash_command="exit 1", on_failure_callback=task_failure_callback ) dag = my_mixed_dag()
注意:如果同时配置了DAG级回调和Task级回调,Task失败时会先触发Task级回调,DAG最终判定为失败时再触发DAG级回调,二者互不冲突。
内容的提问来源于stack exchange,提问作者Chromey
相关产品推荐
相关产品推荐

