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

Airflow DAG任务被跳过:如何实现跳过回调邮件通知?

问题

我想知道是否有类似on_skip_callback的机制来发送通知,目前我用html_email_generator实现邮件通知。我的场景里,某次DAG运行包含4个任务,但每次运行到第三个任务时都会被跳过。我已经配置了failure_callback来处理失败通知,请问需要添加什么才能实现跳过任务的邮件通知?

相关代码

def failire_callback():
     if task_name is 't1':
         print(" OUTPUT : t1 is failed")
     elif task_name is 't2':
         Print(" OUTPUT : t2 is failed")
     elif task_name is 't3':
         print(" OUTPUT : t3 is failed")
     else 
         Print("OUTPUT : t4 is failed")
     subject = Dag has failed
     message = generate_html_template
     send_email(to=IDs, subject = subject, html_content = message)


default_args = (
      Provide_context = True
      on_failure_callback = failure_callback
      dag = dag
)
解决方案

Airflow 提供了on_skip_callback参数,专门用于处理任务被跳过的回调通知,用法和你已配置的on_failure_callback完全一致。

1. 编写跳过任务的回调函数

参照你现有的失败回调逻辑,创建一个处理跳过场景的函数,注意要接收context参数来获取任务信息:

def skip_callback(context):
    task_name = context['task_instance'].task_id
    print(f"OUTPUT : {task_name} is skipped")
    
    subject = "任务已被跳过"
    message = generate_html_template  # 可根据需求生成对应跳过场景的HTML内容
    send_email(to=IDs, subject=subject, html_content=message)

2. 配置on_skip_callback

将on_skip_callback添加到default_args中,全局应用到所有任务:

default_args = {
      "provide_context": True,
      "on_failure_callback": failure_callback,
      "on_skip_callback": skip_callback,  # 新增跳过回调
      "dag": dag
}

如果只想让特定任务(比如你的t3)触发跳过通知,也可以在该任务定义时单独指定on_skip_callback=skip_callback,无需全局配置。

3. 修正原failure_callback的问题

你的原代码存在几处语法和逻辑错误,会导致运行异常,建议一并修正:

  • 函数名拼写错误:failire_callback → failure_callback
  • 缺少context参数,无法获取task_name(依赖provide_context=True传递上下文)
  • 字符串未加引号:subject = Dag has failed → subject = "Dag has failed"
  • 判断字符串相等应使用==而非is
  • Python内置函数print需小写,不要写成Print

修正后的failure_callback示例:

def failure_callback(context):
    task_name = context['task_instance'].task_id
    if task_name == 't1':
        print(" OUTPUT : t1 is failed")
    elif task_name == 't2':
        print(" OUTPUT : t2 is failed")
    elif task_name == 't3':
        print(" OUTPUT : t3 is failed")
    else:
        print("OUTPUT : t4 is failed")
    
    subject = "DAG运行失败"
    message = generate_html_template
    send_email(to=IDs, subject=subject, html_content=message)

内容的提问来源于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 10:17:19