微服务落地最头疼的问题,就是服务之间的通信。同步调用一时爽,流量一上来,链路超时、数据库连接被打满、系统雪崩接踵而至。之前在做后端项目重构时,反复在“服务解耦”和“削峰填谷”这两个环节踩坑,网上关于 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 CLI | v2 版本 | 用于命令行创建资源 |
| Java | JDK 8 及以上 | Spring Boot 集成示例使用 |
| Maven | 3.6 及以上 | 管理项目依赖 |
| IDE | IntelliJ 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.yml3. 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 的消息生命周期:
- 生产者调用
SendMessage把消息写入队列。 - 消费者调用
ReceiveMessage拉取消息(消息进入Invisible 状态,对其它消费者不可见)。 - 消费者处理完成后调用
DeleteMessage删除消息。 - 如果消费者处理失败且没有删除消息,在 Visibility Timeout 超时后,消息重新变回可见状态,可被再次消费。
这个机制保证了消息处理的高可用,但同时也意味着消费者的处理逻辑必须设计为幂等——同一条消息可能被投递多次,重复处理不能产生脏数据。
3.3 SNS 与 SQS 的区别对比
| 对比维度 | SNS | SQS |
|---|---|---|
| 通信模式 | 发布/订阅(一对多) | 点对点(一对一) |
| 消息投递方式 | 主动推送 | 消费者拉取 |
| 消息持久化 | 不保留,推送即结束 | 持久化保存,直到被删除 |
| 典型用途 | 广播事件通知 | 异步任务处理 |
| 消费次数 | 每个订阅者都收到 | 每条消息只被一个消费者处理 |
| 消息顺序 | 不保证 | 标准队列不保证,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 QueueArn4.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-sns和aws-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 Policy5.2 权限错误:AccessDeniedException
这个错误最常见的场景是 IAM 用户权限不足,或者 SQS 队列策略条件限制不匹配。
检查顺序:
- IAM 用户是否具有
sns:Publish、sqs:ReceiveMessage、sqs:DeleteMessage权限。 - SQS 队列策略中的
aws:SourceArn是否与实际的 Topic ARN 完全一致(包括 region 和账号 ID)。 - 如果用了跨账号访问,还需要配置额外的资源策略。
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:ReceiveMessage、sqs:DeleteMessage、sqs: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-service6.6 环境隔离
生产环境和开发环境必须使用不同的 Topic 和 Queue。
- 推荐按环境命名:
order-events-dev、order-events-prod。 - 通过环境变量或配置中心管理 ARN,不要硬编码在代码里。
- 不同环境使用不同的 AWS 账号或在同一账号下使用完整隔离的 VPC 配置更安全。
6.7 消息大小限制
SNS 和 SQS 的单个消息大小限制为 256KB(包括消息属性和消息体)。如果需要传输大对象,不要直接塞进消息体,而是:
- 把文件或者数据存入 Amazon S3。
- 在消息体中携带 S3 对象路径和预签名 URL。
- 消费者从 S3 读取数据后处理。
7. 总结与下一步
本文作为 AWS SNS + SQS 微服务架构系列的第一篇,把核心概念、环境准备、控制台与 CLI 实战、Spring Boot 集成代码、常见问题和工程实践都覆盖到了。掌握了这些,你已经可以独立搭建一套基于 SNS + SQS 的事件驱动微服务示例项目。
几个关键操作值得记住:
- SNS 解决一对多广播,SQS 解决异步削峰。
- Fanout 模式是两者组合的核心架构。
- SQS 队列必须配置访问策略,否则收不到 SNS 推送。
- 消费者代码必须幂等,必须手动删除消息。
- 生产环境必须配置死信队列和监控告警。
下一篇可以从这几个方向继续深入:
- SQS FIFO 队列的严格顺序消费实战。
- SNS 消息过滤策略(按消息属性精确过滤)。
- 结合 Lambda 无服务器消费模式。
- 分布式事务补偿方案与消息最终一致性。
如果你在配置或代码运行中遇到其他问题,欢迎在评论区留言,也可以收藏本文备用。动手把示例代码跑一遍,比囫囵吞枣看十篇文章更有用。