Airflow中回调函数内调用SlackNotifier无响应问题排查
Airflow自定义失败回调中SlackNotifier无响应问题解决
问题说明
- 直接调用
SlackNotifier能正常发送Slack通知,但在自定义on_failure_callback回调里调用时完全没反应 - 回调里的日志输出功能正常,说明回调逻辑本身没问题,只有Slack通知部分失效
- 试过用
send_slack_notification()工具函数,结果一样;SlackWebhookOperator虽然能工作,但改起来麻烦还一堆报错 - 实际需求:只有当DAG不是手动触发的时候,才发送失败通知
原问题代码
def dag_run_info(**context): dag_run = context['dag_run'] run_type = dag_run.run_type return run_type def on_failure_callback(context): run_type = context['task_instance'].xcom_pull(task_ids='get_dag_run_info') if run_type == 'manual': return SlackNotifier(slack_conn_id=SLACK_CONNECTION_ID, text=SLACK_MESSAGE, channel=SLACK_CHANNEL) else: logging.info(f'No Slack notification for failed DAG because run_type:{run_type}')
问题根源
- 没有触发发送动作:你在回调里只是创建了
SlackNotifier的实例并返回,但这个类需要主动调用notify()方法才会实际发送消息,光返回实例不会执行任何操作。 - 逻辑写反:需求是“非手动触发时发通知”,但代码里是
run_type == 'manual'时才处理,和需求完全相反。
修复后的代码
def dag_run_info(**context): dag_run = context['dag_run'] run_type = dag_run.run_type return run_type def on_failure_callback(context): run_type = context['task_instance'].xcom_pull(task_ids='get_dag_run_info') if run_type != 'manual': # 符合需求:非手动触发时发送通知 # 创建通知实例后,必须调用notify()才会发消息 slack_notifier = SlackNotifier( slack_conn_id=SLACK_CONNECTION_ID, text=SLACK_MESSAGE, channel=SLACK_CHANNEL ) slack_notifier.notify(context=context) else: logging.info(f'No Slack notification for failed DAG because run_type:{run_type}')
补充提示
- 如果用
send_slack_notification()工具函数,也要在回调里直接调用它,不能只返回函数对象 - 要是想在通知里加更多失败详情,比如DAG ID、任务ID,可以通过
context变量动态拼接消息内容,示例:slack_message = f"DAG {context['dag_run'].dag_id} 执行失败,任务:{context['task_instance'].task_id}" slack_notifier = SlackNotifier( slack_conn_id=SLACK_CONNECTION_ID, text=slack_message, channel=SLACK_CHANNEL )
内容的提问来源于stack exchange,提问作者Tomáš Tráchly
相关产品推荐
相关产品推荐

