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

