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 2.x版本中,SLA检查已整合到调度器,无需单独启动,但需确保调度器正常运行)airflow scheduler --sla-check - 查看SLA检查日志:
- 1.10.x:查看
sla_checker进程的日志文件 - 2.x:查看调度器日志,搜索
sla关键词,确认是否有执行回调的记录
- 1.10.x:查看
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
相关产品推荐
相关产品推荐

