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

如何将Airflow中某一DAG的失败通知聚合为单封邮件?

实现DAG失败任务的聚合邮件通知

完全可以实现,核心思路是通过任务失败回调记录失败信息,再通过收尾任务聚合所有失败数据并发送邮件,以下是具体实现方案:

核心步骤

  • 为DAG中每个任务配置on_failure_callback,在任务失败时将任务ID、失败时间、日志链接等关键信息存入XCom
  • 添加一个收尾任务,设置trigger_rule=TriggerRule.ALL_DONE(确保所有任务执行完毕后无论结果如何都触发),该任务读取所有失败任务的XCom数据,整理成统一的邮件内容后发送

代码示例

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.email import EmailOperator
from airflow.utils.trigger_rule import TriggerRule
from datetime import datetime

# 记录失败任务信息到XCom
def record_task_failure(context):
    task_instance = context['task_instance']
    failure_info = {
        'task_id': task_instance.task_id,
        'dag_id': task_instance.dag_id,
        'execution_date': str(task_instance.execution_date),
        'log_url': task_instance.log_url
    }
    # 用带任务ID的key推送XCom,避免不同任务数据冲突
    task_instance.xcom_push(key=f"failure_{task_instance.task_id}", value=failure_info)

# 聚合失败信息并生成邮件内容
def aggregate_failure_notification(**context):
    dag_run = context['dag_run']
    ti = context['task_instance']
    # 拉取当前DAG运行中所有以"failure_"开头的XCom数据
    failure_xcoms = ti.xcom_pull(dag_id=dag_run.dag_id, execution_date=dag_run.execution_date, key=None)
    failures = [v for k, v in failure_xcoms.items() if k.startswith('failure_')]
    
    if not failures:
        return  # 无失败任务,直接终止
    
    # 格式化邮件内容
    email_subject = f"DAG {dag_run.dag_id} 运行失败汇总"
    email_body = f"DAG {dag_run.dag_id} 在 {dag_run.execution_date} 运行中,共 {len(failures)} 个任务失败:\n\n"
    for idx, fail in enumerate(failures, 1):
        email_body += f"{idx}. 任务ID: {fail['task_id']}\n"
        email_body += f"   失败时间: {fail['execution_date']}\n"
        email_body += f"   日志链接: {fail['log_url']}\n\n"
    
    # 将邮件内容传递给后续邮件任务
    return {'subject': email_subject, 'body': email_body}

# 定义DAG
with DAG(
    dag_id='large_dag_with_aggregate_notification',
    start_date=datetime(2024, 1, 1),
    schedule_interval='@daily',
    catchup=False
) as dag:

    # 示例业务任务1(模拟失败)
    task1 = PythonOperator(
        task_id='task1',
        python_callable=lambda: 1/0,
        on_failure_callback=record_task_failure
    )

    # 示例业务任务2
    task2 = PythonOperator(
        task_id='task2',
        python_callable=lambda: print('任务执行成功'),
        on_failure_callback=record_task_failure
    )

    # 聚合失败信息
    aggregate_task = PythonOperator(
        task_id='aggregate_failures',
        python_callable=aggregate_failure_notification,
        provide_context=True,
        trigger_rule=TriggerRule.ALL_DONE
    )

    # 发送聚合邮件
    send_email = EmailOperator(
        task_id='send_failure_email',
        to='your_notify_email@example.com',
        subject="{{ task_instance.xcom_pull(task_ids='aggregate_failures')['subject'] }}",
        html_content="{{ task_instance.xcom_pull(task_ids='aggregate_failures')['body'] }}",
        trigger_rule=TriggerRule.ALL_DONE
    )

    # 设置任务依赖:所有业务任务执行完毕后再执行聚合和发邮件
    [task1, task2] >> aggregate_task >> send_email

注意事项

  • XCom默认有48KB存储限制,若需记录大量失败细节(如异常栈),建议改用数据库等外部存储替代XCom
  • 确保Airflow的SMTP配置(smtp_host、smtp_user等)已在airflow.cfg中正确设置
  • 可根据需求扩展邮件内容,比如从context['exception']中提取任务失败的异常信息

内容的提问来源于stack exchange,提问作者Nursultan Imanov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 18:01:06