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
相关产品推荐
相关产品推荐

