AWS SNS+SQS微服务事件驱动架构实战:从概念到代码
2026/9/21 2:53:30 网站建设 项目流程

微服务落地最头疼的问题,就是服务之间的通信。同步调用一时爽,流量一上来,链路超时、数据库连接被打满、系统雪崩接踵而至。之前在做后端项目重构时,反复在“服务解耦”和“削峰填谷”这两个环节踩坑,网上关于 AWS 消息服务的资料大多只讲了单个服务怎么用,很少有把 SNS 和 SQS 组合起来做事件驱动架构的完整教程。这篇文章我们就来系统梳理 AWS SNS 与 SQS 在微服务架构中的核心用法,包含概念对比、架构设计、完整可运行的实战代码和排查经验,新手可以按步骤落地,有基础的开发者也能直接参考排错。

1. 背景与核心概念

1.1 微服务之间为什么需要消息队列

在微服务架构刚兴起的时候,服务之间最常见的通信方式是 HTTP 同步调用。比如订单服务调用库存服务、调用用户服务,一个业务流程串起多个接口。

这种模式在业务规模不大的时候运行得很好,但一旦服务数量增多、流量出现峰值,几个问题会迅速暴露:

  • 耦合过重:订单服务依赖库存服务、支付服务、物流服务的接口,任何一个下游服务出现抖动,上游服务就会跟着失败。
  • 突发流量冲击:秒杀、促销活动期间,瞬间涌入的请求量远超服务处理能力,数据库连接数、线程池全部被打满,最终导致整个系统不可用。
  • 扩展困难:下游服务处理能力的提升需要上游配合调整超时时间、重试次数,牵一发而动全身。
  • 数据一致性问题:分布式环境下,一个业务操作需要同时更新多个服务的数据,无法依靠本地事务保证一致性。

消息队列(Message Queue)的出现,就是为了解决这些问题。它的核心思想是引入一个“中间存储层”,服务之间不再直接通信,而是把消息写入队列,由消费者按自己的节奏去处理。

这样一来,生产者和消费者在时间上、空间上都解耦了。订单服务只需要把“订单创建成功”的消息发出去,不需要关心后续有多少服务要处理这条消息、它们什么时候处理完。

1.2 SNS 和 SQS 分别是什么

AWS 提供了两类托管消息服务:

SNS(Simple Notification Service,简单通知服务)

SNS 是发布/订阅(Pub/Sub)模式的消息服务。一个生产者可以把消息发布到 Topic(主题),多个订阅者可以同时收到这条消息的副本。

订阅者可以是 SQS 队列、Lambda 函数、HTTP/HTTPS 端点、邮件地址、移动端推送等。

SNS 的核心特点是:一对多广播。一条消息发布后,所有订阅了该 Topic 的终端都能接收到。

SQS(Simple Queue Service,简单队列服务)

SQS 是分布式消息队列服务,采用点对点(Point-to-Point)模式。生产者把消息发送到队列,消费者从队列中拉取消息进行处理。

SQS 的核心特点是:一对一消费。一条消息被某个消费者成功接收并删除后,其他消费者不会再次读到这条消息。

1.3 SNS 与 SQS 组合的架构模式

单独使用 SNS 或 SQS 都能解决一部分问题,但在真实微服务架构中,两者经常组合使用,形成Fanout(扇出)架构

生产者 -> SNS Topic -> SQS 队列 A -> 服务 A -> SQS 队列 B -> 服务 B -> SQS 队列 C -> 服务 C

在这个模式下,生产者只往 SNS Topic 发送一条消息,SNS 会自动把这条消息推送到所有订阅的 SQS 队列中。每个队列由不同的微服务独立消费,互不干扰。

这种模式非常适合事件驱动架构:一个业务事件发生后,需要触发多个下游服务的处理动作,但各服务的处理逻辑、处理速度、失败重试策略完全不同。

举个实际例子:用户下单成功后,可能需要同时执行以下操作:

  • 发送订单确认短信。
  • 更新用户积分。
  • 通知仓储系统预留库存。
  • 触发数据分析埋点。

如果四个动作都等待同步完成,下单接口的耗时会非常长。使用 SNS + SQS 后,下单服务只需要往 Topic 发一条消息,四个订阅队列各自去消费,接口可以立即返回。

2. 环境准备与版本说明

在进入实战之前,我们需要先确认环境。下面这些工具和权限是本文示例需要的,你可以根据自己本地的实际情况调整版本。

2.1 运行环境清单

项目建议配置说明
操作系统Linux / macOS / Windows本文命令以 Linux/macOS 为例
AWS 账号必须有需要开通 SNS、SQS 服务
AWS CLIv2 版本用于命令行创建资源
JavaJDK 8 及以上Spring Boot 集成示例使用
Maven3.6 及以上管理项目依赖
IDEIntelliJ IDEA / Eclipse按个人习惯选择

版本说明:AWS 服务处于持续迭代中,SNS、SQS 的 API 相对稳定,但控制台界面和控制台按钮位置可能发生变化。本文重点讲解核心配置思路,如果你的控制台页面显示与截图不同,以实际操作环境的引导为准。

2.2 IAM 权限准备

使用 AWS 服务前,必须先配置 IAM 身份和权限。这里的关键原则是最小权限——只授予实际操作所需的最小权限范围,不要直接使用 AdministratorAccess。

创建一个用于开发的 IAM 用户,建议附加以下内联策略:

{ "Version": "2012-10-17", "Statement": [ { "Effect": "Allow", "Action": [ "sns:CreateTopic", "sns:Subscribe", "sns:Publish", "sns:ListTopics", "sns:ListSubscriptionsByTopic" ], "Resource": "*" }, { "Effect": "Allow", "Action": [ "sqs:CreateQueue", "sqs:GetQueueAttributes", "sqs:SetQueueAttributes", "sqs:SendMessage", "sqs:ReceiveMessage", "sqs:DeleteMessage", "sqs:ListQueues" ], "Resource": "*" } ] }

注意:Resource设置为*仅用于开发环境快速验证。在生产环境中,建议把 Resource 限定到具体的 Topic ARN 和 Queue ARN,例如:

{ "Effect": "Allow", "Action": "sns:Publish", "Resource": "arn:aws:sns:us-east-1:123456789012:order-events" }

2.3 AWS CLI 配置

安装 AWS CLI v2 后,执行配置命令:

aws configure

按提示输入:

  • AWS Access Key ID
  • AWS Secret Access Key
  • Region(例如us-east-1
  • 输出格式(建议填json

验证配置是否生效:

aws sts get-caller-identity

正常输出会包含你的账号 ID、ARN 和 UserId 信息。

2.4 示例项目结构

本文的 Spring Boot 集成示例使用如下结构:

aws-sns-sqs-demo/ ├── pom.xml └── src └── main ├── java │ └── com │ └── example │ └── demo │ ├── DemoApplication.java │ ├── config │ │ └── AwsConfig.java │ ├── model │ │ └── OrderEvent.java │ ├── publisher │ │ └── OrderEventPublisher.java │ └── consumer │ │ └── OrderEventConsumer.java └── resources └── application.yml

3. SNS 与 SQS 核心概念拆解

3.1 SNS 的发布/订阅模型

SNS 的核心组件是 Topic(主题)。Topic 是一个逻辑上的消息通道,生产者向 Topic 发布消息,Topic 负责把消息推送给所有订阅者。

SNS 支持多种订阅终端类型:

  • SQS 队列:把消息推送到 SQS,由消费者拉取处理。
  • Lambda 函数:消息到达时自动触发 Lambda 执行。
  • HTTP/HTTPS 端点:AWS 通过 POST 请求把消息推送到指定 URL。
  • Email / Email-JSON:发送邮件通知。
  • SMS:发送短信通知。
  • 移动端推送:推送到 iOS、Android 等移动应用。

在微服务架构中,最常见的组合是 SNS 推送到 SQS,因为 SQS 提供了消息持久化、批量拉取、延迟队列等能力,比直接 HTTP 推送更可靠。

SNS 的 Topic 属性中,有一个关键概念叫Delivery Policy(重试策略)。当消息推送到某个订阅终端失败时,SNS 会按照策略自动重试。默认重试策略如下:

  • 重试次数:3 次(实际值可能在控制台显示略有差异)
  • 重试间隔:按指数退避增大
  • 如果没有配置死信队列,超过重试次数后消息会丢失

在关键业务场景中,强烈建议为 SNS 订阅配置DLQ(Dead Letter Queue,死信队列)。这样推送失败的消息会进入死信队列,方便后续排查和处理。

3.2 SQS 的队列模型

SQS 提供两种队列类型:

标准队列(Standard Queue)

  • 高吞吐,近乎无限的消息数量。
  • 消息可能乱序,可能重复。
  • 适合对顺序要求不高的场景。

FIFO 队列(First-In-First-Out Queue)

  • 严格保证消息顺序。
  • 消息恰好一次处理(配合去重机制)。
  • 吞吐量限制为每秒 300 次事务(可以批处理提高效率)。
  • 队列名称必须以.fifo结尾。

选择建议:大部分微服务异步处理场景使用标准队列即可,只有在金融交易、库存扣减等对顺序有强依赖的场景才需要 FIFO 队列。

SQS 的消息生命周期:

  1. 生产者调用SendMessage把消息写入队列。
  2. 消费者调用ReceiveMessage拉取消息(消息进入Invisible 状态,对其它消费者不可见)。
  3. 消费者处理完成后调用DeleteMessage删除消息。
  4. 如果消费者处理失败且没有删除消息,在 Visibility Timeout 超时后,消息重新变回可见状态,可被再次消费。

这个机制保证了消息处理的高可用,但同时也意味着消费者的处理逻辑必须设计为幂等——同一条消息可能被投递多次,重复处理不能产生脏数据。

3.3 SNS 与 SQS 的区别对比

对比维度SNSSQS
通信模式发布/订阅(一对多)点对点(一对一)
消息投递方式主动推送消费者拉取
消息持久化不保留,推送即结束持久化保存,直到被删除
典型用途广播事件通知异步任务处理
消费次数每个订阅者都收到每条消息只被一个消费者处理
消息顺序不保证标准队列不保证,FIFO 保证
延迟处理不支持支持 DelaySeconds 延迟队列

简单记忆方式:SNS 负责“把消息发给谁”,SQS 负责“把消息存下来等谁取”

3.4 Fanout 模式(扇出模式)

Fanout 是 SNS + SQS 最经典的组合方式。生产者只关心“事件发生了”,不关心“谁在处理”。SNS Topic 把所有订阅的 SQS 队列都推送一遍,各队列的消费者独立处理。

架构图可以用文字描述为:

+--> SQS Queue A --> Service A SNS Topic ------> SQS Queue B --> Service B +--> SQS Queue C --> Service C

这种模式的优势非常明显:

  • 新增下游服务零成本:新服务只需要创建一个 SQS 队列并订阅该 Topic 即可,上游不需要做任何改动。
  • 故障隔离:某个消费服务宕机,其它队列的消息照常处理。
  • 削峰填谷:流量激增时,消息在队列中堆积,消费方按自身能力慢慢处理。

4. 完整实战案例

接下来我们通过一个完整的订单事件案例,演示如何创建 SNS Topic、SQS 队列、配置订阅,以及如何在 Spring Boot 中集成消息发布和消费。

案例场景:用户下单成功后,订单服务发布一条order.created事件。事件需要同时触发:

  • 通知消费者:订单创建成功,需要发送确认短信。
  • 通知分析服务:记录用户的浏览行为。

所以我们创建一个 SNS Topic,两个 SQS 队列,分别由两个服务消费。

4.1 创建 SNS Topic 和 SQS 队列

先通过 AWS CLI 创建 SNS Topic。

# 创建 SNS Topic aws sns create-topic --name order-events # 查看 Topic 列表 aws sns list-topics

创建成功后,会返回一个 Topic ARN(Amazon Resource Name),类似:

arn:aws:sns:us-east-1:123456789012:order-events

这个 ARN 是 SNS Topic 的唯一标识,后续配置订阅权限时要用到。

接着创建两个 SQS 队列:

# 创建发送短信服务对应的队列 aws sqs create-queue --queue-name order-sms-service # 创建分析服务对应的队列 aws sqs create-queue --queue-name order-analytics-service

创建完成后,可以通过以下命令获取队列的 URL 和 ARN:

# 获取队列 URL aws sqs get-queue-url --queue-name order-sms-service # 获取队列 ARN aws sqs get-queue-attributes --queue-url <QUEUE_URL> --attribute-names QueueArn

4.2 创建订阅关系

把两个 SQS 队列订阅到 SNS Topic 上:

# 订阅 SQS 队列到 SNS Topic aws sns subscribe \ --topic-arn arn:aws:sns:us-east-1:123456789012:order-events \ --protocol sqs \ --notification-endpoint arn:aws:sqs:us-east-1:123456789012:order-sms-service aws sns subscribe \ --topic-arn arn:aws:sns:us-east-1:123456789012:order-events \ --protocol sqs \ --notification-endpoint arn:aws:sqs:us-east-1:123456789012:order-analytics-service

这一步做完,订阅关系已经建立。但有一个关键问题:SQS 默认不允许 SNS 推送消息进来。必须在 SQS 队列的访问策略中显式授权 SNS 发送消息。

4.3 配置 SQS 队列访问策略

修改 SQS 队列的访问策略,允许来自指定 SNS Topic 的消息写入。

aws sqs set-queue-attributes \ --queue-url <QUEUE_URL> \ --attributes '{ "Policy": "{\"Version\":\"2012-10-17\",\"Statement\":[{\"Effect\":\"Allow\",\"Principal\":{\"Service\":\"sns.amazonaws.com\"},\"Action\":\"sqs:SendMessage\",\"Resource\":\"arn:aws:sqs:us-east-1:123456789012:order-sms-service\",\"Condition\":{\"ArnEquals\":{\"aws:SourceArn\":\"arn:aws:sns:us-east-1:123456789012:order-events\"}}}]}" }'

需要给两个队列都配置一遍,把 Queue URL 和 Resource ARN 替换成对应值。

这里解释一下策略的关键点:

  • Principal设置为sns.amazonaws.com,表示允许 SNS 服务访问。
  • Condition限制aws:SourceArn只能是我们指定的 Topic ARN,这是安全加固的重要一步,避免其它 Topic 也能往这个队列写消息。

4.4 发布消息验证 Fanout

现在我们来测试一下,向 SNS Topic 发布一条消息,看看两个队列是否都能收到。

# 发布消息到 SNS Topic aws sns publish \ --topic-arn arn:aws:sns:us-east-1:123456789012:order-events \ --message '{"orderId":"20240601001","userId":"u1001","amount":99.9,"event":"order.created"}'

然后分别从两个队列接收消息:

# 从第一个队列接收消息 aws sqs receive-message --queue-url <QUEUE_URL_1> # 从第二个队列接收消息 aws sqs receive-message --queue-url <QUEUE_URL_2>

如果配置正确,两个队列都应该能读到同一条消息内容。这说明 SNS 的 Fanout 能力生效了。

4.5 Spring Boot 集成

命令行验证通过后,我们来看 Java 代码集成。首先在pom.xml中添加依赖:

<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>com.amazonaws</groupId> <artifactId>aws-java-sdk-sns</artifactId> <version>1.12.600</version> </dependency> <dependency> <groupId>com.amazonaws</groupId> <artifactId>aws-java-sdk-sqs</artifactId> <version>1.12.600</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <scope>test</scope> </dependency> </dependencies>

版本说明:aws-java-sdk-snsaws-java-sdk-sqs的版本会持续更新,上述版本号仅为示例。实际使用时建议到 AWS SDK for Java 官方文档查看最新版本,或者使用 Maven 依赖管理工具自动解析。

然后配置application.yml

aws: region: us-east-1 sns: topic-arn: arn:aws:sns:us-east-1:123456789012:order-events sqs: sms-queue-url: https://sqs.us-east-1.amazonaws.com/123456789012/order-sms-service analytics-queue-url: https://sqs.us-east-1.amazonaws.com/123456789012/order-analytics-service

如果本地没有配置 AWS CLI 的默认凭证,可以通过环境变量传入 Access Key 和 Secret Key:

export AWS_ACCESS_KEY_ID=your_access_key export AWS_SECRET_ACCESS_KEY=your_secret_key

创建 AWS SNS 客户端配置类:

package com.example.demo.config; import com.amazonaws.auth.DefaultAWSCredentialsProviderChain; import com.amazonaws.regions.Regions; import com.amazonaws.services.sns.AmazonSNS; import com.amazonaws.services.sns.AmazonSNSClientBuilder; import com.amazonaws.services.sqs.AmazonSQS; import com.amazonaws.services.sqs.AmazonSQSClientBuilder; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class AwsConfig { @Value("${aws.region}") private String region; @Bean public AmazonSNS amazonSNS() { return AmazonSNSClientBuilder.standard() .withRegion(region) .withCredentials(new DefaultAWSCredentialsProviderChain()) .build(); } @Bean public AmazonSQS amazonSQS() { return AmazonSQSClientBuilder.standard() .withRegion(region) .withCredentials(new DefaultAWSCredentialsProviderChain()) .build(); } }

定义订单事件模型:

package com.example.demo.model; public class OrderEvent { private String orderId; private String userId; private Double amount; private String event; public OrderEvent() { } public OrderEvent(String orderId, String userId, Double amount, String event) { this.orderId = orderId; this.userId = userId; this.amount = amount; this.event = event; } // getter 和 setter 省略 public String getOrderId() { return orderId; } public void setOrderId(String orderId) { this.orderId = orderId; } public String getUserId() { return userId; } public void setUserId(String userId) { this.userId = userId; } public Double getAmount() { return amount; } public void setAmount(Double amount) { this.amount = amount; } public String getEvent() { return event; } public void setEvent(String event) { this.event = event; } @Override public String toString() { return "OrderEvent{" + "orderId='" + orderId + '\'' + ", userId='" + userId + '\'' + ", amount=" + amount + ", event='" + event + '\'' + '}'; } }

编写消息发布者:

package com.example.demo.publisher; import com.amazonaws.services.sns.AmazonSNS; import com.amazonaws.services.sns.model.PublishRequest; import com.amazonaws.services.sns.model.PublishResult; import com.example.demo.model.OrderEvent; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; @Component public class OrderEventPublisher { private final AmazonSNS amazonSNS; private final ObjectMapper objectMapper; @Value("${aws.sns.topic-arn}") private String topicArn; public OrderEventPublisher(AmazonSNS amazonSNS, ObjectMapper objectMapper) { this.amazonSNS = amazonSNS; this.objectMapper = objectMapper; } public String publishOrderCreatedEvent(OrderEvent event) { try { String message = objectMapper.writeValueAsString(event); PublishRequest publishRequest = new PublishRequest() .withTopicArn(topicArn) .withMessage(message) .withMessageGroupId("order-events"); PublishResult result = amazonSNS.publish(publishRequest); return result.getMessageId(); } catch (JsonProcessingException e) { throw new RuntimeException("订单事件序列化失败", e); } } }

编写订单服务接口,演示发送事件:

package com.example.demo.consumer; import com.amazonaws.services.sqs.AmazonSQS; import com.amazonaws.services.sqs.model.Message; import com.amazonaws.services.sqs.model.ReceiveMessageRequest; import com.amazonaws.services.sqs.model.DeleteMessageRequest; import com.example.demo.model.OrderEvent; import com.fasterxml.jackson.databind.ObjectMapper; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import javax.annotation.PostConstruct; import java.util.List; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; @Component public class OrderEventConsumer { private static final Logger logger = LoggerFactory.getLogger(OrderEventConsumer.class); private final AmazonSQS amazonSQS; private final ObjectMapper objectMapper; @Value("${aws.sqs.sms-queue-url}") private String smsQueueUrl; @Value("${aws.sqs.analytics-queue-url}") private String analyticsQueueUrl; private final ScheduledExecutorService executorService = Executors.newScheduledThreadPool(2); public OrderEventConsumer(AmazonSQS amazonSQS, ObjectMapper objectMapper) { this.amazonSQS = amazonSQS; this.objectMapper = objectMapper; } @PostConstruct public void startConsumers() { executorService.scheduleWithFixedDelay(() -> consumeMessages(smsQueueUrl, "SMS_SERVICE"), 0, 5, TimeUnit.SECONDS); executorService.scheduleWithFixedDelay(() -> consumeMessages(analyticsQueueUrl, "ANALYTICS_SERVICE"), 0, 5, TimeUnit.SECONDS); } private void consumeMessages(String queueUrl, String serviceName) { ReceiveMessageRequest receiveMessageRequest = new ReceiveMessageRequest() .withQueueUrl(queueUrl) .withMaxNumberOfMessages(10) .withWaitTimeSeconds(5); List<Message> messages = amazonSQS.receiveMessage(receiveMessageRequest).getMessages(); for (Message message : messages) { try { OrderEvent event = objectMapper.readValue(message.getBody(), OrderEvent.class); logger.info("[{}] 收到订单事件: {}", serviceName, event); // 模拟不同的业务处理逻辑 if ("SMS_SERVICE".equals(serviceName)) { sendSmsNotification(event); } else { recordAnalytics(event); } // 处理成功后删除消息 amazonSQS.deleteMessage(new DeleteMessageRequest() .withQueueUrl(queueUrl) .withReceiptHandle(message.getReceiptHandle())); } catch (Exception e) { logger.error("[{}] 处理消息失败,消息ID: {}", serviceName, message.getMessageId(), e); } } } private void sendSmsNotification(OrderEvent event) { logger.info("模拟发送短信: 用户 {} 的订单 {} 已创建,金额 {}", event.getUserId(), event.getOrderId(), event.getAmount()); } private void recordAnalytics(OrderEvent event) { logger.info("模拟记录分析数据: 订单 {} 关联用户 {}", event.getOrderId(), event.getUserId()); } }

启动类:

package com.example.demo; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; @SpringBootApplication public class DemoApplication { public static void main(String[] args) { SpringApplication.run(DemoApplication.class, args); } }

4.6 运行与验证

现在启动 Spring Boot 应用。应用启动后,你可以通过任意 HTTP 调用触发事件发布,也可以写一个简单的 CommandLineRunner 测试:

package com.example.demo; import com.example.demo.model.OrderEvent; import com.example.demo.publisher.OrderEventPublisher; import org.springframework.boot.CommandLineRunner; import org.springframework.stereotype.Component; @Component public class EventPublishRunner implements CommandLineRunner { private final OrderEventPublisher publisher; public EventPublishRunner(OrderEventPublisher publisher) { this.publisher = publisher; } @Override public void run(String... args) { OrderEvent event = new OrderEvent( "20240601001", "u1001", 99.9, "order.created" ); String messageId = publisher.publishOrderCreatedEvent(event); System.out.println("发布成功,MessageId: " + messageId); } }

运行后,日志中会出现两个消费者的输出,分别模拟发送短信和记录分析数据。

这说明订单事件已经被 Fanout 到两个队列,并被独立消费。

5. 常见问题与排查思路

在使用 SNS + SQS 的过程中,最容易遇到下面几类问题。

5.1 消息发送成功但队列收不到

问题现象常见原因解决思路
SNS 发布返回成功,但 SQS 队列为空SQS 队列策略未授权 SNS检查队列访问策略,确认sqs:SendMessage的权限和Condition设置正确
订阅存在但收不到消息订阅状态不是Confirmed使用控制台检查订阅状态,必要时删除重建
消费者收不到消息消费者使用的队列 URL 错误确认application.yml中的队列 URL 与 CLI 查到的 URL 一致
消息处理失败被反复投递消费者没有删除消息处理成功必须调用DeleteMessage

排查命令参考:

# 查看订阅列表和状态 aws sns list-subscriptions-by-topic --topic-arn <TOPIC_ARN> # 查看队列属性中的策略 aws sqs get-queue-attributes --queue-url <QUEUE_URL> --attribute-names Policy

5.2 权限错误:AccessDeniedException

这个错误最常见的场景是 IAM 用户权限不足,或者 SQS 队列策略条件限制不匹配。

检查顺序:

  1. IAM 用户是否具有sns:Publishsqs:ReceiveMessagesqs:DeleteMessage权限。
  2. SQS 队列策略中的aws:SourceArn是否与实际的 Topic ARN 完全一致(包括 region 和账号 ID)。
  3. 如果用了跨账号访问,还需要配置额外的资源策略。

5.3 消息重复消费

标准队列本身是at-least-once模型,即至少一次投递,极端情况下可能重复。这是 SQS 的设计特性,不是 bug。

解决重复消费的方式:

  • 消费逻辑设计为幂等。
  • 使用消息中的业务唯一键(如orderId)在数据库或 Redis 中做去重。
  • 关键业务场景改用 FIFO 队列。

5.4 SQS 消息 Invisible 时间设置不当

如果消费者处理时间超过 Visibility Timeout,消息会被重新投递,可能造成重复消费。

建议设置:

aws sqs set-queue-attributes \ --queue-url <QUEUE_URL> \ --attributes '{"VisibilityTimeout": "120"}'

或者在生产代码中动态调整:

ReceiveMessageRequest request = new ReceiveMessageRequest() .withQueueUrl(queueUrl) .withVisibilityTimeout(120) .withMaxNumberOfMessages(10) .withWaitTimeSeconds(5);

这里需要根据业务处理时长合理评估,设置太短会重复消费,太长会阻塞消息释放。

5.5 死信队列没有配置

没有配置死信队列时,消息处理失败达到最大接收次数后会被直接丢弃。这种消息丢失对关键业务来说是致命的。

建议为每个 SQS 队列配置死信队列:

# 创建死信队列 aws sqs create-queue --queue-name order-sms-service-dlq # 把原队列的 redrive policy 指向死信队列 aws sqs set-queue-attributes \ --queue-url <QUEUE_URL> \ --attributes '{ "RedrivePolicy": "{\"deadLetterTargetArn\":\"arn:aws:sqs:us-east-1:123456789012:order-sms-service-dlq\",\"maxReceiveCount\":5}" }'

当消息被接收超过 5 次仍未成功处理时,会自动进入死信队列,方便之后单独分析和修复。

6. 最佳实践与工程建议

把 SNS + SQS 真正用好,不是简单地把服务串起来。下面这些实践经验来自真实项目落地过程中的总结。

6.1 IAM 权限最小化

生产环境不要使用Resource: "*"的粗粒度策略。建议对每个 Topic 和每个 Queue 单独配置 ARN 级别的权限,并遵循以下原则:

  • 生产者只给sns:Publish权限。
  • 消费者只给sqs:ReceiveMessagesqs:DeleteMessagesqs:GetQueueAttributes权限。
  • 队列策略中必须用Condition限定aws:SourceArn,防止其他 Topic 恶意投递。

6.2 消息结构统一规范

消息体的设计直接影响后续维护成本。建议定义统一的消息封装结构:

{ "eventId": "uuid", "eventType": "order.created", "version": "1.0", "timestamp": "2024-06-01T12:00:00Z", "payload": { "orderId": "20240601001", "userId": "u1001" } }
  • eventId用于幂等和追踪。
  • eventType用于消费者判断处理逻辑。
  • version便于消息结构升级。
  • payload存放业务数据。

这样设计后,无论消费端如何演进,都能保持相对稳定的解析逻辑。

6.3 消费者必须幂等

这条原则再怎么强调都不过分。SQS 标准队列是 at-least-once 模型,消息几乎一定会重复投递。消费者的处理逻辑应该做到:

  • 数据库写操作使用唯一约束。
  • 先查询后更新,确保操作可重入。
  • 使用 Redis SetNX 或数据库锁做去重。

6.4 合理设置 Visibility Timeout

Visibility Timeout 是 SQS 防重复消费的第一道防线。判断标准是:

  • 消费者处理一条消息的平均耗时是多少。
  • 耗时波动范围有多大。

建议把 Visibility Timeout 设置为平均处理耗时的 6 倍以上,留足安全余量。

6.5 监控和告警

消息堆积、死信队列有消息,这些都是微服务架构中需要重点关注的信号。

建议在 CloudWatch 中配置以下指标告警:

  • ApproximateNumberOfMessagesVisible超过阈值。
  • NumberOfMessagesReceived长时间为零。
  • 死信队列消息数大于零。

用命令查看队列监控指标:

aws cloudwatch list-metrics --namespace AWS/SQS --dimensions Name=QueueName,Value=order-sms-service

6.6 环境隔离

生产环境和开发环境必须使用不同的 Topic 和 Queue。

  • 推荐按环境命名:order-events-devorder-events-prod
  • 通过环境变量或配置中心管理 ARN,不要硬编码在代码里。
  • 不同环境使用不同的 AWS 账号或在同一账号下使用完整隔离的 VPC 配置更安全。

6.7 消息大小限制

SNS 和 SQS 的单个消息大小限制为 256KB(包括消息属性和消息体)。如果需要传输大对象,不要直接塞进消息体,而是:

  1. 把文件或者数据存入 Amazon S3。
  2. 在消息体中携带 S3 对象路径和预签名 URL。
  3. 消费者从 S3 读取数据后处理。

7. 总结与下一步

本文作为 AWS SNS + SQS 微服务架构系列的第一篇,把核心概念、环境准备、控制台与 CLI 实战、Spring Boot 集成代码、常见问题和工程实践都覆盖到了。掌握了这些,你已经可以独立搭建一套基于 SNS + SQS 的事件驱动微服务示例项目。

几个关键操作值得记住:

  • SNS 解决一对多广播,SQS 解决异步削峰。
  • Fanout 模式是两者组合的核心架构。
  • SQS 队列必须配置访问策略,否则收不到 SNS 推送。
  • 消费者代码必须幂等,必须手动删除消息。
  • 生产环境必须配置死信队列和监控告警。

下一篇可以从这几个方向继续深入:

  • SQS FIFO 队列的严格顺序消费实战。
  • SNS 消息过滤策略(按消息属性精确过滤)。
  • 结合 Lambda 无服务器消费模式。
  • 分布式事务补偿方案与消息最终一致性。

如果你在配置或代码运行中遇到其他问题,欢迎在评论区留言,也可以收藏本文备用。动手把示例代码跑一遍,比囫囵吞枣看十篇文章更有用。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询