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

使用boto3订阅Lambda到SNS时触发器未生效报权限错误

问题根因
  • 控制台手动创建SNS到Lambda的订阅时,AWS后端会自动完成两项隐式操作:一是为目标Lambda添加资源级权限策略,授权sns.amazonaws.com服务主体对该函数执行lambda:InvokeFunction操作;二是同步建立双向关联,在Lambda控制台生成对应的SNS触发器条目。
  • 现有boto3脚本仅调用SNS的subscribe接口创建了SNS侧的订阅记录,完全缺失了为Lambda添加调用权限的步骤,这是触发AccessDenied报错、Lambda控制台不显示触发器的核心原因:SNS侧虽有订阅,但Lambda侧既没有授权SNS调用,也没有可识别的关联SNS源信息,自然无法同步触发器、也不允许SNS发起调用。
  • 代码存在一处非核心配置问题:自定义SNS主题访问策略中Resource字段为空,不符合IAM策略规范,后续可能引发其他权限异常。
  • 额外注意:如果SNS主题使用自定义KMS密钥加密,还需要给sns.amazonaws.com服务主体授予密钥的kms:Decrypt、kms:GenerateDataKey权限,否则会出现消息投递失败,使用AWS托管的aws/sns默认密钥无需额外配置。
代码修正方案

修正核心逻辑是在创建SNS订阅前,先调用Lambda的add_permission接口为目标函数添加最小权限的调用授权,限定仅当前创建的SNS主题可调用该Lambda,避免权限过宽;同时补全SNS策略的Resource字段。

关键修正点

  1. 新增Lambda boto3客户端初始化
  2. 创建SNS主题后先补全带正确Resource字段的访问策略,再配置其他主题属性
  3. 创建SNS订阅前,为Lambda添加关联SNS主题的调用权限,增加StatementId去重逻辑避免重复报错
  4. 统一通过主题ARN自动解析账号、区域信息,减少硬编码配置错误

修正后可运行代码片段

import boto3
import json
import logging
from botocore.exceptions import ClientError

logger = logging.getLogger(__name__)
boto3.setup_default_session(profile_name='xxx')


class SnsWrapper:
    """Encapsulates Amazon SNS topic and subscription functions."""

    def __init__(self, sns_resource):
        """
        :param sns_resource: A Boto3 Amazon SNS resource.
        """
        self.sns_resource = sns_resource

    def create_topic(self, name, attributes=None):
        """
        Creates a notification topic.

        :param name: The name of the topic to create.
        :param attributes: Attributes of topic.
        :return: The newly created topic.
        """
        try:
            if attributes:
                topic = self.sns_resource.create_topic(Name=name, Attributes=attributes)
            else:
                topic = self.sns_resource.create_topic(Name=name)
            logger.info("Created topic %s with ARN %s.", name, topic.arn)
        except ClientError:
            logger.exception("Couldn't create topic %s.", name)
            raise
        else:
            return topic

    def list_topics(self):
        """
        Lists topics for the current account.

        :return: An iterator that yields the topics.
        """
        try:
            topics_iter = self.sns_resource.topics.all()
            logger.info("Got topics.")
        except ClientError:
            logger.exception("Couldn't get topics.")
            raise
        else:
            return topics_iter

    @staticmethod
    def delete_topic(topic):
        """
        Deletes a topic. All subscriptions to the topic are also deleted.
        """
        try:
            topic.delete()
            logger.info("Deleted topic %s.", topic.arn)
        except ClientError:
            logger.exception("Couldn't delete topic %s.", topic.arn)
            raise

    @staticmethod
    def subscribe(topic, protocol, endpoint):
        """
        Subscribes an endpoint to the topic. Some endpoint types, such as email,
        must be confirmed before their subscriptions are active. When a subscription
       is not confirmed, its Amazon Resource Number (ARN) is set to
       'PendingConfirmation'.

        :param topic: The topic to subscribe to.
        :param protocol: The protocol of the endpoint, such as 'sms' or 'email'.
        :param endpoint: The endpoint that receives messages, such as a phone number
                     (in E.164 format) for SMS messages, or an email address for
                     email messages.
        :return: The newly added subscription.
        """
        try:
            subscription = topic.subscribe(
                Protocol=protocol, Endpoint=endpoint, ReturnSubscriptionArn=True)
            logger.info("Subscribed %s %s to topic %s.", protocol, endpoint, topic.arn)
        except ClientError:
            logger.exception(
                "Couldn't subscribe %s %s to topic %s.", protocol, endpoint, topic.arn)
            raise
        else:
            return subscription

    def list_subscriptions(self, topic=None):
        """
        Lists subscriptions for the current account, optionally limited to a
        specific topic.

        :param topic: When specified, only subscriptions to this topic are returned.
        :return: An iterator that yields the subscriptions.
        """
        try:
            if topic is None:
                subs_iter = self.sns_resource.subscriptions.all()
            else:
                subs_iter = topic.subscriptions.all()
            logger.info("Got subscriptions.")
        except ClientError:
            logger.exception("Couldn't get subscriptions.")
            raise
        else:
            return subs_iter

    @staticmethod
    def add_subscription_filter(subscription, attributes):
        """
        Adds a filter policy to a subscription. A filter policy is a key and a
        list of values that are allowed. When a message is published, it must have an
        attribute that passes the filter or it will not be sent to the subscription.

        :param subscription: The subscription the filter policy is attached to.
        :param attributes: A dictionary of key-value pairs that define the filter.
        """
        try:
            att_policy = {key: [value] for key, value in attributes.items()}
            subscription.set_attributes(
                AttributeName='FilterPolicy', AttributeValue=json.dumps(att_policy))
            logger.info("Added filter to subscription %s.", subscription.arn)
        except ClientError:
            logger.exception(
                "Couldn't add filter to subscription %s.", subscription.arn)
            raise

    @staticmethod
    def delete_subscription(subscription):
        """
        Unsubscribes and deletes a subscription.
        """
        try:
            subscription.delete()
            logger.info("Deleted subscription %s.", subscription.arn)
        except ClientError:
            logger.exception("Couldn't delete subscription %s.", subscription.arn)
            raise

    def publish_text_message(self, phone_number, message):
        """
        Publishes a text message directly to a phone number without need for a
        subscription.

        :param phone_number: The phone number that receives the message. This must be
                         in E.164 format. For example, a United States phone
                         number might be +12065550101.
        :param message: The message to send.
        :return: The ID of the message.
        """
        try:
            response = self.sns_resource.meta.client.publish(
                PhoneNumber=phone_number, Message=message)
            message_id = response['MessageId']
            logger.info("Published message to %s.", phone_number)
        except ClientError:
            logger.exception("Couldn't publish message to %s.", phone_number)
            raise
        else:
            return message_id

    @staticmethod
    def publish_message(topic, message, attributes):
        """
        Publishes a message, with attributes, to a topic. Subscriptions can be filtered
        based on message attributes so that a subscription receives messages only
        when specified attributes are present.

        :param topic: The topic to publish to.
        :param message: The message to publish.
        :param attributes: The key-value attributes to attach to the message. Values
                       must be either `str` or `bytes`.
        :return: The ID of the message.
        """
        try:
            att_dict = {}
            for key, value in attributes.items():
                if isinstance(value, str):
                    att_dict[key] = {'DataType': 'String', 'StringValue': value}
                elif isinstance(value, bytes):
                    att_dict[key] = {'DataType': 'Binary', 'BinaryValue': value}
            response = topic.publish(Message=message, MessageAttributes=att_dict)
            message_id = response['MessageId']
            logger.info(
                "Published message with attributes %s to topic %s.", attributes,
                topic.arn)
        except ClientError:
            logger.exception("Couldn't publish message to topic %s.", topic.arn)
            raise
        else:
            return message_id

    @staticmethod
    def publish_multi_message(
            topic, subject, default_message, sms_message, email_message):
        """
        Publishes a multi-format message to a topic. A multi-format message takes
        different forms based on the protocol of the subscriber. For example,
        an SMS subscriber might receive a short, text-only version of the message
        while an email subscriber could receive an HTML version of the message.

        :param topic: The topic to publish to.
        :param subject: The subject of the message.
        :param default_message: The default version of the message. This version is
                            sent to subscribers that have protocols that are not
                            otherwise specified in the structured message.
        :param sms_message: The version of the message sent to SMS subscribers.
        :param email_message: The version of the message sent to email subscribers.
        :return: The ID of the message.
        """
        try:
            message = {
                'default': default_message,
                'sms': sms_message,
                'email': email_message
            }
            response = topic.publish(
                Message=json.dumps(message), Subject=subject, MessageStructure='json')
            message_id = response['MessageId']
            logger.info("Published multi-format message to topic %s.", topic.arn)
        except ClientError:
            logger.exception("Couldn't publish message to topic %s.", topic.arn)
            raise
        else:
            return message_id


# -----TEST SNS-----
logging.basicConfig(level=logging.INFO, format='%(levelname)s: %(message)s')
sns = boto3.resource('sns')
# 新增Lambda客户端初始化
lambda_client = boto3.client('lambda')

sns_wrapper = SnsWrapper(sns)

# 先创建基础主题,拿到ARN后再配置全量属性
topic = sns_wrapper.create_topic(name='test-sns-functionality')
topic_arn = topic.arn
account_id = topic_arn.split(':')[4]
region = topic_arn.split(':')[3]

# 组装投递策略
retry_policy = json.dumps({
    "http": {
        "defaultHealthyRetryPolicy": {
            "minDelayTarget": 20,
            "maxDelayTarget": 21,
            "numRetries": 3,
            "numMaxDelayRetries": 0,
            "numNoDelayRetries": 0,
            "numMinDelayRetries": 0,
            "backoffFunction": "linear"
        },
        "disableSubscriptionOverrides": False
    }
})

# 组装正确的SNS访问策略,补全Resource字段
example_access_policy = json.dumps({
    "Version": "2008-10-17",
    "Id": "__default_policy_ID",
    "Statement": [
        {
            "Sid": "__default_statement_ID",
            "Effect": "Allow",
            "Principal": {
                "AWS": "*"
            },
            "Action": [
                "SNS:Publish",
                "SNS:RemovePermission",
                "SNS:SetTopicAttributes",
                "SNS:DeleteTopic",
                "SNS:ListSubscriptionsByTopic",
                "SNS:GetTopicAttributes",
                "SNS:AddPermission",
                "SNS:Subscribe"
            ],
            "Resource": topic_arn,
            "Condition": {
                "StringEquals": {
                    "AWS:SourceOwner": account_id
                }
            }
        }
    ]
})

# 更新主题全量属性
kms_key_arn = f'arn:aws:kms:{region}:{account_id}:key/a0a18770-6a62-418c-b35f-0adee2616270'
topic.set_attributes(AttributeName='DeliveryPolicy', AttributeValue=retry_policy)
topic.set_attributes(AttributeName='Policy', AttributeValue=example_access_policy)
topic.set_attributes(AttributeName='KmsMasterKeyId', AttributeValue=kms_key_arn)

lambda_arn = f'arn:aws:lambda:{region}:{account_id}:function:sns-test-lambda'

# 核心新增:为Lambda添加SNS调用权限
try:
    lambda_client.add_permission(
        FunctionName=lambda_arn,
        StatementId=f'sns-invoke-{topic.name}',
        Action='lambda:InvokeFunction',
        Principal='sns.amazonaws.com',
        SourceArn=topic_arn,
        SourceAccount=account_id
    )
    logger.info(f"Added Lambda invoke permission for SNS topic {topic_arn}")
except ClientError as e:
    if e.response['Error']['Code'] == 'ResourceConflictException':
        logger.info("Lambda permission already exists, skip adding")
    else:
        raise

# 创建SNS到Lambda的订阅
subscriber_arn = sns_wrapper.subscribe(topic, 'lambda', lambda_arn)

# 发布测试消息
sns_wrapper.publish_message(topic, f'This is message from topic: {topic_arn}', attributes={'atr1': 'val1'})
验证说明
  • 代码执行完成后等待1-2分钟,刷新Lambda控制台即可看到对应SNS触发器条目
  • 发布测试消息不会再抛出AccessDenied异常,Lambda可正常接收消息并触发执行
  • 如果删除SNS订阅,需要同步调用Lambda的remove_permission接口删除对应权限语句,避免残留无效权限

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.02 06:06:48