Airflow 2.5.1中DagFactory生成DAG的回调未触发,Slack通知失败求助
核心问题梳理
使用DagFactory通过YAML配置Airflow DAG时,配置了任务/ DAG的启动、成功、失败Slack回调,但仅日志显示默认回调记录,Slack未收到通知。以下是针对性排查点及修正方案:
一、YAML配置的明显错误
参数拼写错误
retry_dely应为retry_delay,参数名错误会导致默认参数加载异常,间接影响回调逻辑生效。文件路径错误
on_success_callback_file中使用//开头路径,改为/home/airflow/provider/dags/dependencies/slack_notification.py,重复斜线会导致DagFactory无法定位回调文件。回调名称不匹配
on_execute_callback_name配置为slack_notification_start_operator,但Python代码中对应的函数名是slack_notification_start,名称完全不匹配会导致DagFactory无法找到目标回调函数,直接跳过执行。
二、Python代码的兼容性与逻辑问题
Airflow 2.x模块迁移
Airflow 2.x中airflow.contrib.operators.slack_webhook_operator已迁移至airflow.providers.slack.operators.slack_webhook,使用旧模块会导致导入失败,回调函数静默报错。回调函数参数缺失
slack_notification_failed_callback未添加**context参数,Airflow回调传参时会触发参数不匹配错误,导致函数执行终止。Slack Operator参数冲突
SlackWebhookOperator同时传入http_conn_id和webhook_token,若http_conn_id对应的连接已配置完整的Slack Webhook URL,重复传webhook_token会导致请求参数冲突,通知发送失败。缺失调试日志
函数内未添加日志打印,无法判断是回调未执行还是Slack请求失败。建议在每个回调函数开头添加日志记录。
三、环境与权限问题
文件权限检查
确认slack_notification.py的权限设置,Airflow Worker用户需具备文件读取权限,可在Worker节点执行ls -l /home/airflow/provider/dags/dependencies/slack_notification.py验证。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已包含验证信息
- Conn Type选择
回调层级混淆
on_execute_callback是任务级回调(每个任务执行前触发),若需DAG启动时发送通知,应使用on_dag_run_callback(DAG级回调),需在YAML的dag节点下单独配置。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

