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
相关产品推荐
相关产品推荐

