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

Airflow DAG失败后无法触发AWS SNS主题,求排查解决

问题描述

我创建了一个测试用Airflow DAG,运行时会主动触发失败,但DAG执行失败后无法触发对应的AWS SNS主题。以下是我编写的用于在DAG失败时触发SNS主题的Python代码,请问是代码存在问题,还是代码之外有遗漏的配置导致功能无法正常工作?

from datetime import datetime
from airflow.models import DAG
from airflow.operators.python_operator import PythonOperator
from airflow.operators.dummy_operator import DummyOperator
from pathlib import Path
from airflow.providers.amazon.aws.operators.sns import SnsPublishOperator
import boto3
from airflow.providers.amazon.aws.hooks.sns import SnsHook

DAG_NAME = 'dag_name'

# Funciton which triggers SNS notification on failure
def failure_callback(context):
    sns_notification = SnsPublishOperator(
        sns_hook = SnsHook(sns_topic_arn='My Topic ARN'),
        dag_id = context["dag"].dag_id,
        task_id = context["task_instance"].task_id,
        exception = str(context.get("exception")),
        log_url = context["task_instance"].log_url
    )

# Trigger SNS notification
    sns_client = boto3.client('sns', region_name='eu-west-2')
    sns_client.publish(
        TopicArn='My Topic ARN',
        Message=f"DAG {context['dag'].dag_id} failed on task {context['task_instance'].task_id}"
    )

# Function to intentionally fail DAG (ONLY FOR TESTING)
def fail_dag():
    raise Exception("Intentional DAG failure for test")

# Define DAG
default_args = {
    "on_failure_callback": failure_callback,
    'owner': 'Me',
    'start_date': datetime(2021, 12, 9),
    'is_prod': False
}

with DAG(
    dag_id='dag_name',
    default_args=default_args,
    schedule_interval=None,
    catchup=False,
) as dag:
    
    # To fail task intentionally
    fail_task = PythonOperator(
        task_id='fail_task',
        python_callable=fail_dag,
    )

    # Dummy start and end tasks for DAG
    start_task = DummyOperator(task_id='start')
    end_task = DummyOperator(task_id='end')

    # Dependencies (chained)
    start_task >> fail_task >> end_task

globals()[DAG_NAME] = dag
问题排查与解决

一、代码层面的问题

  • 冗余代码未执行:你在failure_callback里初始化了SnsPublishOperator,但从未调用它的execute方法,这段代码完全无效,直接删除即可。
  • boto3客户端未复用Airflow配置:直接用boto3.client创建SNS客户端,不会使用Airflow中配置的AWS连接信息,容易出现身份验证问题。建议改用Airflow官方的SnsHook来获取客户端,复用Airflow的AWS连接配置,修改后的回调函数如下:
def failure_callback(context):
    # 替换为你在Airflow中配置的AWS连接ID
    sns_hook = SnsHook(aws_conn_id='aws_default')
    sns_client = sns_hook.get_client()
    sns_client.publish(
        TopicArn='My Topic ARN',
        Message=f"DAG {context['dag'].dag_id} failed on task {context['task_instance'].task_id}"
    )
  • 缩进错误:原代码中# Trigger SNS notification下方的代码缩进层级错误,虽然Python可能未抛出语法错误,但逻辑上必须属于failure_callback函数内部,需确保缩进正确。

二、配置层面的可能遗漏

  • Airflow AWS连接配置:必须在Airflow中创建有效的AWS连接(可通过UI或airflow connections add命令)。如果使用IAM角色(如EC2实例角色、EKS Pod角色),需确保Airflow运行环境拥有该角色权限;如果使用访问密钥,需确保密钥具备sns:Publish权限。
  • SNS主题权限策略:检查SNS主题的访问策略,允许Airflow运行环境的IAM身份(用户/角色)执行sns:Publish动作。示例策略如下:
{
    "Effect": "Allow",
    "Principal": {
        "AWS": "arn:aws:iam::你的AWS账号ID:user/airflow-service-user"
    },
    "Action": "sns:Publish",
    "Resource": "你的SNS主题ARN"
}
  • 提供商包兼容性:确保apache-airflow-providers-amazon包的版本与你的Airflow版本兼容(Airflow 2.x需搭配对应版本的提供商包),避免版本不兼容导致Hook/Operator失效。
  • 回调触发范围:default_args中的on_failure_callback仅在任务失败时触发,若DAG未启动(如调度配置错误),该回调不会执行。你的测试场景中fail_task失败会触发回调,需确认任务确实进入失败状态。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 13:02:02