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

Airflow 2.5.1中DagFactory生成DAG的回调未触发,Slack通知失败求助

Airflow 2.5.1 + DagFactory Slack回调未触发排查方案

核心问题梳理

使用DagFactory通过YAML配置Airflow DAG时,配置了任务/ DAG的启动、成功、失败Slack回调,但仅日志显示默认回调记录,Slack未收到通知。以下是针对性排查点及修正方案:


一、YAML配置的明显错误

  1. 参数拼写错误
    retry_dely应为retry_delay,参数名错误会导致默认参数加载异常,间接影响回调逻辑生效。

  2. 文件路径错误
    on_success_callback_file中使用//开头路径,改为/home/airflow/provider/dags/dependencies/slack_notification.py,重复斜线会导致DagFactory无法定位回调文件。

  3. 回调名称不匹配
    on_execute_callback_name配置为slack_notification_start_operator,但Python代码中对应的函数名是slack_notification_start,名称完全不匹配会导致DagFactory无法找到目标回调函数,直接跳过执行。


二、Python代码的兼容性与逻辑问题

  1. Airflow 2.x模块迁移
    Airflow 2.x中airflow.contrib.operators.slack_webhook_operator已迁移至airflow.providers.slack.operators.slack_webhook,使用旧模块会导致导入失败,回调函数静默报错。

  2. 回调函数参数缺失
    slack_notification_failed_callback未添加**context参数,Airflow回调传参时会触发参数不匹配错误,导致函数执行终止。

  3. Slack Operator参数冲突
    SlackWebhookOperator同时传入http_conn_id和webhook_token,若http_conn_id对应的连接已配置完整的Slack Webhook URL,重复传webhook_token会导致请求参数冲突,通知发送失败。

  4. 缺失调试日志
    函数内未添加日志打印,无法判断是回调未执行还是Slack请求失败。建议在每个回调函数开头添加日志记录。


三、环境与权限问题

  1. 文件权限检查
    确认slack_notification.py的权限设置,Airflow Worker用户需具备文件读取权限,可在Worker节点执行ls -l /home/airflow/provider/dags/dependencies/slack_notification.py验证。

  2. Slack连接配置验证
    在Airflow UI的Connections中,确认my_slack_connection_id的配置:

    • Conn Type选择HTTP
    • Host填写完整的Slack Webhook URL(如https://hooks.slack.com/services/XXX/XXX/XXX)
    • 无需填写Password,Webhook URL已包含验证信息
  3. 回调层级混淆
    on_execute_callback是任务级回调(每个任务执行前触发),若需DAG启动时发送通知,应使用on_dag_run_callback(DAG级回调),需在YAML的dag节点下单独配置。

  4. DagFactory版本兼容性
    旧版本DagFactory对Airflow 2.x的回调支持有限,建议升级至兼容Airflow 2.5.1的DagFactory版本。


修正后的示例代码

修正后的YAML配置

dag:
  default_args:
    owner: "me"
    depends_on_past: False
    retries: 0
    retry_delay: datetime.timedelta(minutes=3)
    # 任务级执行前回调
    on_execute_callback_name: slack_notification_start
    on_execute_callback_file: /home/airflow/provider/dags/dependencies/slack_notification.py
    # 任务级成功回调
    on_success_callback_name: slack_notification_finish
    on_success_callback_file: /home/airflow/provider/dags/dependencies/slack_notification.py
    # 任务级失败回调
    on_failure_callback_name: slack_notification_failed_callback
    on_failure_callback_file: /home/airflow/provider/dags/dependencies/slack_notification.py
    provide_context: True
    wait_for_downstream: True
  # 可选:DAG级启动回调
  on_dag_run_callback_name: slack_notification_dag_start
  on_dag_run_callback_file: /home/airflow/provider/dags/dependencies/slack_notification.py

修正后的Python代码

import logging
from datetime import datetime

from airflow.hooks.base import BaseHook
from airflow.providers.slack.operators.slack_webhook import SlackWebhookOperator

SLACK_CONN_ID = "my_slack_connection_id"
logger = logging.getLogger(__name__)

def slack_notification_start(**context):
    logger.info("开始执行Slack任务启动通知")
    slack_op = SlackWebhookOperator(
        task_id='slack_notification_start',
        http_conn_id=SLACK_CONN_ID,
        message=f"""
        :stopwatch: Airflow任务启动
        *任务ID*: {context['task_instance'].task_id}
        *Dag*: {context['task_instance'].dag_id}
        *执行时间*: {context['execution_date']}
        *日志链接*: {context['task_instance'].log_url}
        """,
        username='airflow'
    )
    return slack_op.execute(context=context)

def slack_notification_failed_callback(**context):
    logger.info("开始执行Slack任务失败通知")
    current_time = datetime.now()
    slack_op = SlackWebhookOperator(
        task_id='slack_notification_fail',
        http_conn_id=SLACK_CONN_ID,
        message=f"""
        :red_circle: 任务执行失败
        *任务*: {context['task_instance'].task_id}
        *Dag*: {context['task_instance'].dag_id}
        *失败时间*: {current_time}
        *日志链接*: {context['task_instance'].log_url}
        """,
        username='airflow'
    )
    return slack_op.execute(context=context)

def slack_notification_finish(**context):
    logger.info("开始执行Slack任务完成通知")
    current_time = datetime.now()
    slack_op = SlackWebhookOperator(
        task_id='slack_notification_finish',
        http_conn_id=SLACK_CONN_ID,
        message=f"""
        :white_check_mark: Airflow任务完成
        *Dag*: {context['task_instance'].dag_id}
        *完成时间*: {current_time}
        *日志链接*: {context['task_instance'].log_url}
        """,
        username='airflow'
    )
    return slack_op.execute(context=context)

# DAG级启动回调示例
def slack_notification_dag_start(**context):
    logger.info("开始执行Slack DAG启动通知")
    slack_op = SlackWebhookOperator(
        task_id='slack_notification_dag_start',
        http_conn_id=SLACK_CONN_ID,
        message=f"""
        :rocket: Airflow DAG启动
        *Dag*: {context['dag_run'].dag_id}
        *执行时间*: {context['dag_run'].execution_date}
        """,
        username='airflow'
    )
    return slack_op.execute(context=context)

内容的提问来源于stack exchange,提问作者Prof. Falken

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 19:05:23