如何从Airflow失败任务实例中获取报错原因/异常信息?
获取Airflow失败任务实例的异常原因方法
针对你遇到的TaskInstance.task为None无法获取异常的问题,可通过以下几种方式直接获取失败任务的异常信息:
直接读取TaskInstance的内置异常属性
Airflow的TaskInstance对象本身存储了异常相关数据,无需依赖task字段:from airflow.utils.state import TaskInstanceState failed_tis = dag_run.get_task_instances([TaskInstanceState.FAILED]) for ti in failed_tis: # 获取异常对象,转成字符串查看详情 exception = ti.exception if exception: print(f"任务 {ti.task_id} 失败原因: {str(exception)}") # 获取错误描述文本 error_msg = ti.error if error_msg: print(f"错误描述: {error_msg}")动态映射任务的实例会自动带上索引后缀(如
source1__map_index_0),task_id会准确对应到具体失败的映射实例。通过XCom在错误回调中传递异常信息
既然on_error_callback已用于自定义逻辑,可在回调中新增XCom推送异常的逻辑,后续在sink任务中拉取:- 修改错误回调函数:
def custom_error_callback(context): ti = context["ti"] # 推送异常信息到XCom,key可自定义 ti.xcom_push(key="failure_detail", value=str(context["exception"])) # 原有自定义逻辑代码... - 在sink任务中拉取:
failed_tis = dag_run.get_task_instances([TaskInstanceState.FAILED]) for ti in failed_tis: # 拉取对应任务实例的异常信息 exception_detail = ti.xcom_pull(key="failure_detail", task_ids=ti.task_id) if exception_detail: print(f"任务 {ti.task_id} 失败详情: {exception_detail}")
- 修改错误回调函数:
查询Airflow日志模型获取完整异常栈
如果需要更完整的异常堆栈信息,可通过Airflow的ORM模型直接查询日志记录:from airflow.models import Log failed_tis = dag_run.get_task_instances([TaskInstanceState.FAILED]) for ti in failed_tis: logs = Log.query.filter( Log.dag_id == ti.dag_id, Log.task_id == ti.task_id, Log.execution_date == ti.execution_date, Log.try_number == ti.try_number ).all() # 提取日志中的异常内容(通常日志末尾会有堆栈信息) for log in logs: if "Traceback" in log.message: print(f"任务 {ti.task_id} 异常堆栈:\n{log.message}")
以上方法均无需依赖TaskInstance.task字段,适配动态映射任务场景,且不会与已有的on_error_callback或Teardown任务约束冲突。
内容的提问来源于stack exchange,提问作者Vito De Tullio
相关产品推荐
相关产品推荐

