如何将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

