You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Airflow 2.0.2中on_retry_callback异常无日志问题及解决问询

Airflow 2.0.2 on_retry_callback异常日志问题解答

这是默认行为吗?

是的,在Airflow 2.0.2版本中,on_retry_callback回调函数抛出的异常不会自动记录到task_instance日志中,属于默认行为。Airflow早期版本的回调机制设计中,不会主动捕获并关联回调异常到任务实例日志,这类异常通常只会在调度器日志中短暂出现,或者直接被静默处理,导致调试回调逻辑时很难定位问题。

如何开启该异常的日志记录?

方法1:在回调函数内手动添加日志捕获

这是最直接且可靠的方案,在回调逻辑中显式捕获异常并写入日志,同时保留完整的堆栈信息:

import logging
from airflow.models import TaskInstance

def cleanup_on_retry(context):
    ti: TaskInstance = context["ti"]
    # 获取与任务实例绑定的日志对象
    logger = logging.getLogger(f"airflow.task.task_instance.{ti.dag_id}.{ti.task_id}")
    try:
        # 你的清理操作逻辑
        print("执行清理步骤...")
        # 示例:模拟抛出异常
        raise ValueError("清理操作失败,测试异常日志")
    except Exception as e:
        # 记录异常到任务实例日志,exc_info=True会输出完整堆栈
        logger.error(f"重试回调清理失败: {str(e)}", exc_info=True)
        # 可选:如果需要中断重试流程,可重新抛出异常;否则注释此行
        raise

使用与任务实例绑定的日志对象,能确保异常信息直接写入对应任务的task_instance日志文件中,方便后续调试。

方法2:临时排查可查看调度器日志

如果不想修改代码,可先查看Airflow调度器的日志文件(默认路径为$AIRFLOW_HOME/logs/scheduler),回调函数的异常信息可能会被调度器日志捕获,但这种方式无法将异常与具体任务实例日志关联,仅适合临时排查。

注意事项

  • 若在回调中重新抛出异常,会中断当前的任务重试流程;如果希望任务继续重试,可去掉raise语句,仅记录日志即可。
  • Airflow 2.2+版本对回调日志机制有优化,会自动捕获回调异常并记录到任务实例日志,但2.0.2版本无此特性,只能通过手动日志记录解决。

内容的提问来源于stack exchange,提问作者Venkatesan Muniappan

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 23:40:21