如何获取Airflow排队任务名称以设置超时告警?
默认的airflow.executor.queued_tasks是聚合类指标,仅统计排队任务总数,不携带任务名称、DAG ID等维度信息。要实现需求,需要通过自定义指标采集或主动查询任务状态两种方式补充维度数据,再配置Datadog告警。
一、自定义Airflow指标,添加任务名称维度
通过修改Airflow的Executor逻辑,在任务进入排队状态时,上报带task_id、dag_id维度的自定义指标,同时支持后续计算排队时长。
操作步骤:
自定义Executor代码(以CeleryExecutor为例):
在Airflow的CeleryExecutor类的submit方法中,添加自定义Metrics上报逻辑,利用Airflow内置的Statsd客户端发送带维度的指标:from airflow.executors.celery_executor import CeleryExecutor from airflow.utils.timezone import utcnow from datadog import statsd class CustomCeleryExecutor(CeleryExecutor): def submit(self, key, command, queue=None, executor_config=None): # 解析任务唯一标识中的DAG和任务ID dag_id, task_id, _, _, _ = key # 上报带维度的排队任务标记 statsd.increment( "airflow.queued_task", tags=[f"dag_id:{dag_id}", f"task_id:{task_id}"] ) # 调用父类方法完成任务提交 super().submit(key, command, queue, executor_config)然后在Airflow配置文件
airflow.cfg中指定自定义Executor:executor = path.to.your.CustomCeleryExecutor配置Datadog采集:
确保Airflow的Statsd客户端指向Datadog Agent的Statsd端口(默认8125),Datadog会自动采集带维度的airflow.queued_task指标,后续可通过维度筛选特定任务。
二、通过Airflow元数据库查询,推送任务超时事件到Datadog
如果不想修改Airflow核心代码,可通过定时查询Airflow元数据库,获取排队任务的名称和时长,再推送到Datadog作为事件或自定义指标。
操作步骤:
编写查询脚本(Python示例,适配Postgres元数据库):
连接数据库筛选超时排队任务,计算时长后推送到Datadog:import psycopg2 from datetime import datetime, timedelta from datadog import initialize, api # Datadog初始化(建议从环境变量读取密钥) options = { 'api_key': '<YOUR_DATADOG_API_KEY>', 'app_key': '<YOUR_DATADOG_APP_KEY>' } initialize(**options) # 数据库连接配置 db_config = { 'dbname': 'airflow', 'user': 'airflow', 'password': '<DB_PASSWORD>', 'host': '<DB_HOST>' } def check_queued_tasks(timeout_minutes=30): conn = psycopg2.connect(**db_config) cursor = conn.cursor() now = datetime.utcnow() timeout_threshold = now - timedelta(minutes=timeout_minutes) # 查询超时排队任务 query = """ SELECT dag_id, task_id, queued_at FROM task_instance WHERE state = 'queued' AND queued_at < %s """ cursor.execute(query, (timeout_threshold,)) tasks = cursor.fetchall() for dag_id, task_id, queued_at in tasks: duration = (now - queued_at).total_seconds() / 60 # 发送Datadog告警事件 api.Event.create( title=f"Airflow任务排队超时", text=f"DAG: {dag_id} | 任务: {task_id} | 排队时长: {round(duration, 2)}分钟", tags=[f"dag_id:{dag_id}", f"task_id:{task_id}", "alert:queued_timeout"] ) # 可选:发送自定义时长指标 api.Metric.send( metric='airflow.queued_task.duration', points=(now.timestamp(), duration), tags=[f"dag_id:{dag_id}", f"task_id:{task_id}"] ) cursor.close() conn.close() if __name__ == "__main__": check_queued_tasks(timeout_minutes=30)定时执行脚本:
用Cron或Airflow自身的定时DAG定期运行该脚本(比如每5分钟一次),及时发现超时任务。
三、配置Datadog告警
根据上述方案的输出,配置对应的告警规则:
- 基于自定义指标:在Datadog中创建监控,选择
airflow.queued_task.duration指标,按dag_id、task_id维度筛选,设置阈值为X分钟,触发告警通知。 - 基于事件:创建事件监控,筛选标签为
alert:queued_timeout的事件,当事件出现时触发告警。
内容的提问来源于stack exchange,提问作者Even

