如何在Airflow中实现任务运行超10小时强制失败并告警?
Airflow任务超时强制失败、告警并终止的实现方案
方案一:单任务级别的精准控制(推荐)
直接利用Airflow内置参数结合自定义回调,就能实现需求:
- 设置任务超时时间
在Operator中配置execution_timeout参数,指定10小时的超时阈值。当任务运行时长超过该值时,Airflow会自动发送终止信号(SIGTERM)结束任务进程,并将任务标记为failed。 - 绑定失败告警回调
编写自定义Python函数处理告警逻辑(如发送邮件、企业微信/钉钉消息),通过on_failure_callback参数绑定到任务上,任务失败时自动触发告警。
代码示例:
from datetime import timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.utils.dates import days_ago import requests # 自定义失败告警逻辑 def send_timeout_alert(context): alert_content = ( f"任务超时失败通知\n" f"DAG ID: {context['dag'].dag_id}\n" f"任务ID: {context['task_instance'].task_id}\n" f"执行时间: {context['execution_date']}" ) # 替换为你的告警接口(如钉钉机器人、企业微信API) requests.post( url="https://your-alert-api.com/send", json={"content": alert_content} ) default_args = { 'owner': 'airflow', 'execution_timeout': timedelta(hours=10), # 10小时超时设置 'on_failure_callback': send_timeout_alert, # 绑定告警回调 } with DAG( dag_id='timeout_monitor_dag', default_args=default_args, schedule_interval='@daily', start_date=days_ago(1), catchup=False, ) as dag: def simulate_long_running_task(): # 模拟超过10小时的任务 import time time.sleep(36001) long_task = PythonOperator( task_id='long_running_demo', python_callable=simulate_long_running_task, )
方案二:全局任务统一管控
如果需要给所有任务设置默认超时规则,可以修改Airflow配置文件airflow.cfg:
- 设置
core.task_timeout = 36000(单位:秒,对应10小时),作为全局默认任务超时时间(优先级低于单任务的execution_timeout) - 在全局
default_args中统一配置on_failure_callback,实现所有任务失败时的告警触发
注意事项
- 确保Airflow调度器与工作节点时间同步,避免超时判断出现偏差
- 若使用KubernetesExecutor,
execution_timeout会自动触发K8s Pod的终止逻辑 - 告警函数建议加入重试机制,避免因网络波动导致告警丢失
内容的提问来源于stack exchange,提问作者Sreelakshmi Siri
相关产品推荐
相关产品推荐

