Airflow中on_failure_callback配置Telegram失败告警问题求助
Airflow任务失败时发送Telegram通知的问题
我有一个可用的TelegramOperator函数:
send_telegram_message = TelegramOperator( task_id="send_telegram_message", token="<token>", chat_id="<chat ID>", text="Some notification" )
按t1 >> send_telegram_message的方式串联任务时,该函数可正常工作。但我希望它能在任务失败时发送通知,尝试多种方法均未成功:
将其放入
default_args中,导致DAG报错(UI无具体错误信息):default_args={ "depends_on_past": False, "email_on_failure": False, "email_on_retry": False, "on_failure_callback": send_telegram_message, "retries": 1, "retry_delay": timedelta(seconds=30) },将
on_failure_callback参数放入任务参数中,无论配置on_failure_callback还是on_success_callback均无效果(UI显示任务成功完成,无错误日志):t1 = BashOperator( task_id="entering_virtual_environment", depends_on_past=False, bash_command="source /home/fitwist/airflow/airflow_env/bin/activate", retries=2, on_failure_callback=send_telegram_message )创建高阶函数并配置到DAG中,未收到Telegram消息:
def task_success_callback(context): op = TelegramOperator( task_id="send_telegram_message", token="<token>", chat_id="<chat ID>", text="Some notification" ) op.execute(context) with DAG( "monthly_effectiveness", default_args={ ... }, ... on_failure_callback=send_telegram_message, ... ) as dag: ...在任务流中直接串联,无任何反应(无报错):
t1 >> t2 >> send_telegram_message send_telegram_message >> send_failed_telegram_message
请问如何实现Airflow任务失败时的Telegram通知?
解决方案
核心问题分析
on_failure_callback要求传入可调用函数,而非Operator实例。直接传入TelegramOperator对象会导致Airflow无法解析执行逻辑,这是前三种方法失效的主要原因。
正确实现方式
方法1:自定义回调函数(推荐)
编写接收context参数的函数,内部实例化TelegramOperator并执行,同时可通过context提取任务失败的详细信息,让通知内容更精准:
from airflow.providers.telegram.operators.telegram import TelegramOperator from airflow.utils.context import Context def send_failure_alert(context: Context): # 从上下文提取失败任务信息 task_instance = context.get("task_instance") error_msg = context.get("exception") or "未知错误" notification_text = ( f"任务失败提醒\n" f"DAG ID: {task_instance.dag_id}\n" f"任务ID: {task_instance.task_id}\n" f"失败原因: {str(error_msg)}" ) # 实例化并执行Telegram通知 alert_op = TelegramOperator( task_id=f"telegram_failure_alert_{task_instance.task_id}", token="<你的Telegram Bot Token>", chat_id="<你的Chat ID>", text=notification_text ) alert_op.execute(context=context)
将函数配置到全局或单个任务:
# 全局配置:所有任务失败触发通知 default_args={ "depends_on_past": False, "email_on_failure": False, "email_on_retry": False, "on_failure_callback": send_failure_alert, "retries": 1, "retry_delay": timedelta(seconds=30) } # 单个任务配置:仅该任务失败触发通知 t1 = BashOperator( task_id="entering_virtual_environment", bash_command="source /home/fitwist/airflow/airflow_env/bin/activate", retries=2, on_failure_callback=send_failure_alert )
方法2:利用TriggerRule实现失败分支
通过任务依赖+触发规则,仅当前面任务失败时执行通知任务,适合需要可视化监控通知逻辑的场景:
from airflow.utils.trigger_rule import TriggerRule # 定义正常任务 t1 = BashOperator( task_id="entering_virtual_environment", bash_command="source /home/fitwist/airflow/airflow_env/bin/activate", retries=2 ) # 定义失败通知任务 send_failure_telegram = TelegramOperator( task_id="send_failure_telegram", token="<你的Telegram Bot Token>", chat_id="<你的Chat ID>", text="任务t1执行失败" ) # 设置触发规则:仅t1失败时执行通知 t1 >> send_failure_telegram send_failure_telegram.trigger_rule = TriggerRule.ONE_FAILED
注意事项
- 确保Airflow环境已安装Telegram provider:
pip install apache-airflow-providers-telegram - 验证Telegram Bot的Token和Chat ID正确性,Bot需加入目标聊天组并拥有发消息权限。
- 使用回调函数时,确保正确处理
context参数,避免隐性报错。
内容的提问来源于stack exchange,提问作者Helen Kapatsa
相关产品推荐
相关产品推荐

