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

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

注意事项

  1. 确保Airflow环境已安装Telegram provider:
    pip install apache-airflow-providers-telegram
    
  2. 验证Telegram Bot的Token和Chat ID正确性,Bot需加入目标聊天组并拥有发消息权限。
  3. 使用回调函数时,确保正确处理context参数,避免隐性报错。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 18:43:13