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

Airflow中如何在DAG失败回调中汇总所有任务失败信息

Airflow DAG失败邮件汇总所有失败任务信息的解决方案

核心思路

直接查询Airflow元数据数据库中当前DAG运行实例(DAGRun)下的所有失败任务实例(TaskInstance),在DAG级别的失败回调中汇总这些信息后发送邮件,避开全局变量、XCOM的局限性。

具体实现步骤

1. 编写自定义DAG失败回调函数

这个函数负责查询失败任务、汇总信息并发送邮件,同时可保留原有的Slack告警逻辑:

from airflow.models import DAGRun, TaskInstance
from airflow.utils.email import send_email
from airflow.settings import Session
from datetime import datetime

def custom_dag_failure_callback(context):
    dag_run = context["dag_run"]
    session = Session()

    # 查询当前DAGRun下所有失败的任务实例
    failed_tasks = (
        session.query(TaskInstance)
        .filter(
            TaskInstance.dag_id == dag_run.dag_id,
            TaskInstance.run_id == dag_run.run_id,
            TaskInstance.state == "failed",
        )
        .all()
    )

    # 格式化汇总失败信息
    failure_summary = f"## DAG {dag_run.dag_id} 运行失败汇总\n\n"
    failure_summary += f"DAG运行ID: {dag_run.run_id}\n"
    failure_summary += f"开始时间: {dag_run.start_date.strftime('%Y-%m-%d %H:%M:%S')}\n"
    failure_summary += f"触发告警时间: {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}\n\n"
    failure_summary += "### 失败任务详情:\n"

    for task in failed_tasks:
        failure_summary += f"- **任务ID**: {task.task_id}\n"
        failure_summary += f"  失败时间: {task.end_date.strftime('%Y-%m-%d %H:%M:%S')}\n"
        failure_summary += f"  异常堆栈信息:\n{task.exception}\n"
        failure_summary += f"  日志链接: {task.log_url}\n\n"

    # 发送汇总邮件
    send_email(
        to=["your-notification-email@example.com"],
        subject=f"[Airflow 告警] DAG {dag_run.dag_id} 运行失败",
        html_content=f"<pre style='white-space: pre-wrap;'>{failure_summary}</pre>",
    )

    # 保留原有的Slack告警逻辑
    # dag_failure_notification_alert(context)

2. 配置DAG使用自定义回调

在DAG定义中指定on_failure_callback为上面的自定义函数,同时保留任务级别的Slack告警:

from airflow import DAG
from airflow.operators.postgres_operator import PostgresOperator
from airflow.operators.dummy import DummyOperator
from datetime import datetime, timedelta

default_args = {
    "owner": "airflow",
    "depends_on_past": False,
    "start_date": datetime(2024, 1, 1),
    "email_on_failure": False,  # 关闭默认邮件告警,使用自定义逻辑
    "email_on_retry": False,
    "retries": 0,
}

with DAG(
    "test_alert_notification",
    default_args=default_args,
    description="Test DAG with parallel Postgres tasks and failure alerts",
    schedule_interval=timedelta(days=1),
    catchup=False,
    on_failure_callback=custom_dag_failure_callback,  # 绑定自定义失败回调
) as dag:
    start_task = DummyOperator(task_id="start_task")

    # 创建5个并行Postgres任务
    postgres_tasks = []
    for i in range(5):
        task = PostgresOperator(
            task_id=f"postgres_task_{i}",
            postgres_conn_id="your_postgres_connection_id",
            sql="SELECT 1;  -- 替换为实际SQL",
            on_failure_callback=task_failure_notification_alert,  # 保留任务级Slack告警
        )
        postgres_tasks.append(task)

    end_task = DummyOperator(task_id="end_task")

    start_task >> postgres_tasks >> end_task

为什么之前的方案无效?

  • 全局变量:跨DAGRun共享,并发场景下会被覆盖,无法精准关联当前失败的DAG实例
  • XCOM:仅能传递成功任务的输出,失败任务可能无法推送XCOM;且DAG回调中批量获取所有失败任务的XCOM操作繁琐
  • 普通Python回调:若未结合元数据查询,仅能拿到触发回调的单个任务信息,无法汇总全局失败情况

注意事项

  • 确保Airflow元数据库连接正常,Session能正常执行查询
  • Airflow 2.x版本中导入路径无变化,若使用1.x需调整部分导入(如send_email在airflow.operators.email中)
  • 可根据需求优化邮件格式(如改用HTML表格),提升可读性
  • 若需同时发送Slack告警,直接在自定义回调中调用原有的dag_failure_notification_alert函数即可

内容的提问来源于stack exchange,提问作者Shubhank Gupta

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 08:45:30