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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 02:42:49