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

Airflow技术问题:如何获取被调用失败DAG的日志错误信息

提取Airflow被调用DAG的失败任务错误信息

针对你通过airflow.api.client.local_client.Client.trigger_dag调用其他DAG后,需要在调用方日志中展示被调用DAG错误信息的需求,以下是两种实用的实现方式:

方法一:直接获取任务异常核心信息

利用TaskInstance对象自带的exception和traceback属性,快速获取失败任务的错误类型、消息及回溯信息,适合快速定位问题:

from airflow.exceptions import AirflowException

dag_runs = DagRun.find(dag_id=dag_id, run_id=dri, state="failed")
if len(dag_runs) > 0:
    for dag_run in dag_runs:
        task_instances = dag_run.get_task_instances()
        failed_task_instance = next((ti for ti in task_instances if ti.state == "failed"), None)
        if failed_task_instance:
            error_details = []
            # 提取异常类型和消息
            if failed_task_instance.exception:
                error_details.append(f"错误类型: {type(failed_task_instance.exception).__name__}")
                error_details.append(f"错误消息: {str(failed_task_instance.exception)}")
            # 提取回溯信息
            if failed_task_instance.traceback:
                error_details.append(f"回溯信息:\n{failed_task_instance.traceback}")
            
            error_msg = "\n".join(error_details)
            # 打印到调用方日志
            print(f"被调用DAG [{dag_id}] 的失败任务 [{failed_task_instance.task_id}] 错误详情:\n{error_msg}")
            # 抛出带详细信息的异常
            raise AirflowException(f"{dag_id} DAG执行失败!\n{error_msg}")

方法二:读取任务完整执行日志

如果需要查看任务执行的全过程日志,可使用Airflow的TaskLogReader读取完整日志内容:

from airflow.exceptions import AirflowException
from airflow.utils.log.logging_mixin import LoggingMixin
from airflow.utils.log.log_reader import TaskLogReader

dag_runs = DagRun.find(dag_id=dag_id, run_id=dri, state="failed")
if len(dag_runs) > 0:
    for dag_run in dag_runs:
        task_instances = dag_run.get_task_instances()
        failed_task_instance = next((ti for ti in task_instances if ti.state == "failed"), None)
        if failed_task_instance:
            # 初始化日志读取器
            log_reader = TaskLogReader(LoggingMixin().logger)
            # 读取任务最后100行日志(可根据需求调整start_line和end_line)
            log_lines = log_reader.read_log(
                ti=failed_task_instance,
                try_number=failed_task_instance.try_number,
                start_line=0,
                end_line=100
            )
            full_log = "\n".join(log_lines)
            
            # 打印完整日志到调用方日志
            print(f"被调用DAG [{dag_id}] 的失败任务 [{failed_task_instance.task_id}] 完整执行日志:\n{full_log}")
            # 抛出带日志的异常
            raise AirflowException(f"{dag_id} DAG执行失败!\n完整日志:\n{full_log}")

注意事项

  • exception和traceback仅存储最近一次失败的核心错误,数据量小,适合快速排查
  • 读取完整日志可获取更多执行上下文,但注意控制读取行数避免日志过载
  • 确保调用方任务的执行账号拥有被调用DAG任务日志的访问权限

内容的提问来源于stack exchange,提问作者Shahid Thaika

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 12:00:20