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

如何将Airflow失败任务日志内容附加到故障通知邮件?

问题描述

我有一个运行BashOperator任务的Airflow DAG,任务失败时收到的邮件仅包含简略信息:

Try 1 out of 1
Exception:
Bash command failed. The command returned a non-zero exit code 1.
Log: Link
Host: 2db56ea2ab34
Mark success: Link

我需要获取任务失败的具体原因(即点击Log链接后看到的错误信息),希望Airflow能将该日志直接附加到故障通知邮件中,避免收件人跳转查看。当前使用Airflow 2.6.1版本,必要时可升级。另外,我的BashOperator执行的是docker run命令——这是否是我下述尝试失败的原因?我知晓DockerOperator的存在,但这是另一问题。

我的尝试

我编写了自定义的failure_callback()函数,能成功发送邮件,但无法正确附加日志:

import os
import tempfile
from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from airflow.utils.email import send_email
from airflow.utils.log.logging_mixin import LoggingMixin
from airflow.utils.log.log_reader import TaskLogReader
from datetime import datetime

def failure_callback(context):
    task_instance = context['task_instance']
    log_url = task_instance.log_url
    exception = context.get('exception')

    # Fetch the log content using TaskLogReader
    try:
        log_reader = TaskLogReader()
        log_content, _ = log_reader.read_log_chunks(task_instance, try_number=task_instance.try_number, metadata={})
        log_content = ''.join([chunk['message'] for chunk in log_content[0]])

    except Exception as e:
        log_content = f"Could not fetch log content: {e}"

    subject = f"Airflow Task Failure: {task_instance.task_id}"
    html_content = f"""
    Task: {task_instance.task_id}<br>
    DAG: {task_instance.dag_id}<br>
    Execution Time: {context['logical_date']}<br>
    Log URL: <a href="{log_url}">{log_url}</a><br>
    Exception: {exception}<br>

    """

    # Write log content to a temporary file
    with tempfile.NamedTemporaryFile(delete=False, suffix='.log') as temp_log_file:
        temp_log_file.write(log_content.encode('utf-8'))
        temp_log_file_path = temp_log_file.name

    # Send email with log content as attachment
    send_email(
        to=default_args['email'],
        subject=subject,
        html_content=html_content,
        files=[temp_log_file_path]  # Attach the log file
    )
    os.remove(temp_log_file_path)

已在DAG定义中添加on_failure_callback=failure_callback。

结果

邮件已发送并附带附件,但附件内容显示:

Could not fetch log content: tuple indices must be integers or slices, not str

同时,在Airflow UI中查看任务日志时显示:

[2024-08-08, 14:41:09 EDT] {file_task_handler.py:522} ERROR - Could not read served logs
Traceback (most recent call last):
  File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 1407, in _run_raw_task
    self._execute_task_with_callbacks(context, test_mode)
  File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 1558, in _execute_task_with_callbacks
    result = self._execute_task(context, task_orig)
  File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 1628, in _execute_task
    result = execute_callable(context=context)
  File "/home/airflow/.local/lib/python3.8/site-packages/airflow/operators/bash.py", line 210, in execute
    raise AirflowException(
airflow.exceptions.AirflowException: Bash command failed. The command returned a non-zero exit code 1.
解决方案

1. 修复日志读取逻辑错误

你遇到的tuple indices must be integers or slices, not str错误,是因为Airflow 2.6.x中TaskLogReader.read_log_chunks的返回格式与代码假设不符:该版本返回的log_content是由元组(而非字典)组成的列表,每个元组结构为(日志消息内容, 时间戳),而非包含message键的字典。

修改日志读取部分的代码:

# 原错误代码
log_content = ''.join([chunk['message'] for chunk in log_content[0]])

# 修正后代码
log_content = ''.join([chunk[0] for chunk in log_content[0]])

2. 确保日志读取权限与配置

  • 确认Airflow Worker有权限访问日志存储路径(本地文件系统或S3等远程存储)。
  • 如果使用远程日志存储,需检查airflow.cfg中remote_logging、remote_log_conn_id等配置是否正确,且Worker能正常连接到存储服务。

3. 关于docker run命令的影响

执行docker run命令本身不会导致日志读取失败,但需注意:

  • 若docker run的stderr输出未被捕获,Airflow日志中会缺少具体错误信息,可在命令后添加2>&1将错误输出合并到标准输出,确保Airflow记录完整日志。
  • 此问题与你当前日志读取失败的报错无关。

4. 可选:升级Airflow版本(推荐)

升级到Airflow 2.8+后,可使用更简洁的TaskInstance.log_content属性直接获取日志,无需手动调用TaskLogReader:

try:
    log_content = task_instance.log_content
except Exception as e:
    log_content = f"Could not fetch log content: {e}"

完整修正后的failure_callback函数

import os
import tempfile
from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from airflow.utils.email import send_email
from airflow.utils.log.log_reader import TaskLogReader
from datetime import datetime

def failure_callback(context):
    task_instance = context['task_instance']
    log_url = task_instance.log_url
    exception = context.get('exception')

    # Fetch the log content using TaskLogReader
    try:
        log_reader = TaskLogReader()
        log_content, _ = log_reader.read_log_chunks(task_instance, try_number=task_instance.try_number, metadata={})
        # 从元组中提取日志消息
        log_content = ''.join([chunk[0] for chunk in log_content[0]])

    except Exception as e:
        log_content = f"Could not fetch log content: {e}"

    subject = f"Airflow Task Failure: {task_instance.task_id}"
    html_content = f"""
    Task: {task_instance.task_id}<br>
    DAG: {task_instance.dag_id}<br>
    Execution Time: {context['logical_date']}<br>
    Log URL: <a href="{log_url}">{log_url}</a><br>
    Exception: {exception}<br>
    """

    # Write log content to a temporary file
    with tempfile.NamedTemporaryFile(delete=False, suffix='.log') as temp_log_file:
        temp_log_file.write(log_content.encode('utf-8'))
        temp_log_file_path = temp_log_file.name

    # Send email with log content as attachment
    send_email(
        to=context['default_args']['email'],
        subject=subject,
        html_content=html_content,
        files=[temp_log_file_path]
    )
    os.remove(temp_log_file_path)

额外优化:直接在邮件正文中显示日志

若不想使用附件,可将日志直接嵌入邮件正文(注意日志过长可能影响显示):

html_content = f"""
Task: {task_instance.task_id}<br>
DAG: {task_instance.dag_id}<br>
Execution Time: {context['logical_date']}<br>
Log URL: <a href="{log_url}">{log_url}</a><br>
Exception: {exception}<br><br>
<h3>任务日志:</h3>
<pre>{log_content}</pre>
"""

内容的提问来源于stack exchange,提问作者Sam Firke

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 00:01:07