如何在AWS Airflow中获取跳过的DAG运行及无状态任务邮件告警?
解决Airflow中跳过任务与无状态任务的邮件通知问题
一、处理跳过任务的邮件通知
Airflow原生提供了on_skipped_callback参数,专门用于任务被跳过场景触发回调逻辑,刚好能解决你无法收到跳过任务告警的问题。
1. 编写跳过任务的回调函数
先修正你现有代码的语法错误,再新增跳过任务的回调逻辑(注意要通过context参数获取任务实例信息):
from airflow.utils.email import send_email def failure_callback(context): task_instance = context['task_instance'] task_name = task_instance.task_id print(f"OUTPUT : {task_name} is failed") subject = "DAG任务执行失败" message = generate_html_template(context) # 假设这是你定义的邮件内容生成函数 send_email(to=YOUR_EMAIL_LIST, subject=subject, html_content=message) def skipped_callback(context): task_instance = context['task_instance'] task_name = task_instance.task_id print(f"OUTPUT : {task_name} is skipped") subject = f"DAG任务[{task_name}]被跳过" message = generate_html_template(context) send_email(to=YOUR_EMAIL_LIST, subject=subject, html_content=message)
2. 配置回调到任务
你可以把跳过回调加到default_args里,让所有任务生效:
default_args = { 'provide_context': True, 'on_failure_callback': failure_callback, 'on_skipped_callback': skipped_callback, # 补充你的其他默认参数,比如owner、start_date等 }
如果只需要给第三个任务单独配置,也可以在任务定义时直接指定:
t3 = PythonOperator( task_id='t3', python_callable=your_t3_function, on_skipped_callback=skipped_callback, dag=dag )
二、处理无状态(未执行)任务的通知
第四个任务处于“无状态”,本质是因为上游任务(t3)被跳过,而任务默认触发规则是all_success,导致t4无法被调度执行。要监控这种情况,有两种可行方案:
1. 修改任务触发规则(可选)
如果希望t4即使上游被跳过也执行,可以给t4设置trigger_rule='all_done',这样t4会进入执行状态(成功/失败),就能通过现有回调逻辑触发通知。但如果不需要执行t4,只是要告警,看下面的方案。
2. DAG级监控未执行任务
方案1:DAG运行完成后检查所有任务状态
在DAG的on_success_callback里遍历所有任务实例,排查未执行的任务:
def dag_completion_callback(context): dag_run = context['dag_run'] task_instances = dag_run.get_task_instances() unexecuted_tasks = [ti.task_id for ti in task_instances if ti.state is None] if unexecuted_tasks: subject = f"DAG[{dag_run.dag_id}]存在未执行任务" message = f"以下任务未进入执行状态:{', '.join(unexecuted_tasks)}" send_email(to=YOUR_EMAIL_LIST, subject=subject, html_content=message) # 在DAG定义时配置该回调 dag = DAG( dag_id='your_dag_id', default_args=default_args, on_success_callback=dag_completion_callback, # 补充你的其他DAG参数 )
方案2:添加收尾监控任务
在DAG末尾新增一个监控任务,依赖所有前置任务(设置trigger_rule='all_done'),专门检查t4的状态:
def check_unexecuted_tasks(**context): ti_t4 = context['task_instance'].dag_run.get_task_instance(task_id='t4') if ti_t4.state is None: subject = "任务t4未执行" message = "由于上游任务t3被跳过,t4未进入执行状态" send_email(to=YOUR_EMAIL_LIST, subject=subject, html_content=message) monitor_task = PythonOperator( task_id='monitor_unexecuted_tasks', python_callable=check_unexecuted_tasks, provide_context=True, trigger_rule='all_done', dag=dag ) # 设置依赖关系 t1 >> t2 >> t3 >> t4 >> monitor_task
三、你现有代码的关键错误修正
你提供的代码存在几个必须修正的问题:
- 函数名拼写错误:
failire_callback→failure_callback - Python内置函数大小写错误:
Print→print default_args必须是字典(用{})而非元组(()),键名需小写(provide_context而非Provide_context)- 回调函数必须接收
context参数才能获取任务实例信息 - 字符串需用引号包裹:
subject = "Dag has failed" generate_html_template如果是函数需传入参数并调用:generate_html_template(context)
内容的提问来源于stack exchange,提问作者Gayathri
相关产品推荐
相关产品推荐

