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

Airflow未调用sla_miss_callback函数问题求助

Airflow SLA Slack告警不触发问题的解决思路

问题分析

核心矛盾是:SLA错过记录已正常生成,但sla_miss_callback函数未实际执行——放在DAG定义中时显示notification sent = False;放入default_args中显示notification sent = True,但Slack告警始终未发出。问题主要出在回调函数的实现逻辑、Airflow SLA检查机制的配置两个层面。

具体解决方案

1. 修正SLA回调函数的实现

你之前使用SlackWebhookOperator触发告警的方式,在SLA回调的上下文环境中容易出现初始化异常(比如缺少必要的执行上下文)。建议改用更轻量可靠的SlackWebhookHook直接发送消息:

from airflow.hooks.slack_hook import SlackWebhookHook
from datetime import datetime, timedelta

def slack_failure_notification(context):
    failed_alert = SlackWebhookOperator(
        task_id='slack_failure_notification',
        http_conn_id='datascience_alerts_slack',
        message="failed task")
    return failed_alert.execute(context=context)

def slack_sla_notification(dag, task_list, blocking_task_list, slas, blocking_tis):
    # 拼接SLA错过的详细信息
    sla_details = []
    for sla in slas:
        task_instance = sla.task_instance
        sla_details.append(
            f"Task: {task_instance.task_id}\n"
            f"Execution Date: {task_instance.execution_date}\n"
            f"SLA Deadline: {sla.sla_deadline}\n"
            f"Actual End Time: {task_instance.end_date}"
        )
    message = f"⚠️ Airflow SLA Missed!\n\nAffected Tasks: {task_list}\n\nDetails:\n{chr(10).join(sla_details)}"
    
    # 使用SlackWebhookHook直接发送消息
    slack_hook = SlackWebhookHook(
        http_conn_id='datascience_alerts_slack',
        message=message
    )
    slack_hook.execute()

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2019, 1, 1),
    'depends_on_past': False,
    'retries': 0,
    'on_failure_callback': slack_failure_notification
}

# 必须将sla_miss_callback放在DAG定义中(旧版Airflow不支持在default_args中配置)
dag = DAG(
    'template_dag',
    default_args=default_args,
    catchup=False,
    schedule_interval="10 13 * * *",
    sla_miss_callback=slack_sla_notification
)

test_task = SnowflakeOperator(
    task_id='test_task',
    dag=dag,
    snowflake_conn_id='snowflake-static-datascience_airflow',
    sql="somelongrunningcodehere",
    sla=timedelta(minutes=1)
)

2. 确认Airflow SLA检查机制正常运行

SLA回调由Airflow的SLA检查器触发,而非任务实例本身,需确保以下配置:

  • 检查airflow.cfg中的参数:sla_checker_enabled = True
  • 若使用Airflow 1.10.x版本,需单独启动SLA检查进程:
    airflow scheduler --sla-check
    
    (Airflow 2.x版本中,SLA检查已整合到调度器,无需单独启动,但需确保调度器正常运行)
  • 查看SLA检查日志:
    • 1.10.x:查看sla_checker进程的日志文件
    • 2.x:查看调度器日志,搜索sla关键词,确认是否有执行回调的记录

3. 关键注意事项

  • 禁止在default_args中配置sla_miss_callback:旧版Airflow(1.10.x)不支持该参数放入default_args,即使配置后显示notification sent = True,实际也不会触发回调。
  • 任务实例详情无sla_miss_callback属性是正常现象:该回调是DAG级别的配置,并非任务实例属性,任务只需配置sla参数即可。
  • 验证回调触发:手动触发任务后,故意让任务运行超过1分钟(你的SLA设置),通过airflow sla_misses命令查看记录,确认notification_sent字段变为True,同时检查Slack是否收到消息。

临时替代方案(若上述方法无效)

如果SLA回调机制因版本bug无法正常工作,可采用以下临时方案:
利用任务的on_success_callback或on_failure_callback,手动检查任务执行时间是否超过SLA:

def check_sla_and_notify(context):
    task_instance = context['task_instance']
    sla = task_instance.task.sla
    if sla and (task_instance.end_date - task_instance.start_date) > sla:
        # 发送SLA告警
        slack_hook = SlackWebhookHook(
            http_conn_id='datascience_alerts_slack',
            message=f"Task {task_instance.task_id} exceeded SLA!"
        )
        slack_hook.execute()
    # 原有失败告警逻辑
    if context['task_instance'].state == 'failed':
        slack_failure_notification(context)

# 将该函数配置为任务的on_success_callback和on_failure_callback

内容的提问来源于stack exchange,提问作者Hanmyo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 03:16:11