AWS SNS+SQS微服务事件驱动架构实战:从概念到代码

发布时间:2026/8/27 4:58:49
AWS SNS+SQS微服务事件驱动架构实战:从概念到代码 微服务落地最头疼的问题就是服务之间的通信。同步调用一时爽流量一上来链路超时、数据库连接被打满、系统雪崩接踵而至。之前在做后端项目重构时反复在“服务解耦”和“削峰填谷”这两个环节踩坑网上关于 AWS 消息服务的资料大多只讲了单个服务怎么用很少有把 SNS 和 SQS 组合起来做事件驱动架构的完整教程。这篇文章我们就来系统梳理 AWS SNS 与 SQS 在微服务架构中的核心用法包含概念对比、架构设计、完整可运行的实战代码和排查经验新手可以按步骤落地有基础的开发者也能直接参考排错。1. 背景与核心概念1.1 微服务之间为什么需要消息队列在微服务架构刚兴起的时候服务之间最常见的通信方式是 HTTP 同步调用。比如订单服务调用库存服务、调用用户服务一个业务流程串起多个接口。这种模式在业务规模不大的时候运行得很好但一旦服务数量增多、流量出现峰值几个问题会迅速暴露耦合过重订单服务依赖库存服务、支付服务、物流服务的接口任何一个下游服务出现抖动上游服务就会跟着失败。突发流量冲击秒杀、促销活动期间瞬间涌入的请求量远超服务处理能力数据库连接数、线程池全部被打满最终导致整个系统不可用。扩展困难下游服务处理能力的提升需要上游配合调整超时时间、重试次数牵一发而动全身。数据一致性问题分布式环境下一个业务操作需要同时更新多个服务的数据无法依靠本地事务保证一致性。消息队列Message Queue的出现就是为了解决这些问题。它的核心思想是引入一个“中间存储层”服务之间不再直接通信而是把消息写入队列由消费者按自己的节奏去处理。这样一来生产者和消费者在时间上、空间上都解耦了。订单服务只需要把“订单创建成功”的消息发出去不需要关心后续有多少服务要处理这条消息、它们什么时候处理完。1.2 SNS 和 SQS 分别是什么AWS 提供了两类托管消息服务SNSSimple Notification Service简单通知服务SNS 是发布/订阅Pub/Sub模式的消息服务。一个生产者可以把消息发布到 Topic主题多个订阅者可以同时收到这条消息的副本。订阅者可以是 SQS 队列、Lambda 函数、HTTP/HTTPS 端点、邮件地址、移动端推送等。SNS 的核心特点是一对多广播。一条消息发布后所有订阅了该 Topic 的终端都能接收到。SQSSimple 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 IDAWS Secret Access KeyRegion例如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 订阅配置DLQDead 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 的区别对比对比维度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 ARNAmazon 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 groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdcom.amazonaws/groupId artifactIdaws-java-sdk-sns/artifactId version1.12.600/version /dependency dependency groupIdcom.amazonaws/groupId artifactIdaws-java-sdk-sqs/artifactId version1.12.600/version /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-test/artifactId scopetest/scope /dependency /dependencies版本说明aws-java-sdk-sns和aws-java-sdk-sqs的版本会持续更新上述版本号仅为示例。实际使用时建议到 AWS SDK for Java 官方文档查看最新版本或者使用 Maven 依赖管理工具自动解析。然后配置application.ymlaws: 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 Keyexport AWS_ACCESS_KEY_IDyour_access_key export AWS_SECRET_ACCESS_KEYyour_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); ListMessage 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 TimeoutVisibility Timeout 是 SQS 防重复消费的第一道防线。判断标准是消费者处理一条消息的平均耗时是多少。耗时波动范围有多大。建议把 Visibility Timeout 设置为平均处理耗时的 6 倍以上留足安全余量。6.5 监控和告警消息堆积、死信队列有消息这些都是微服务架构中需要重点关注的信号。建议在 CloudWatch 中配置以下指标告警ApproximateNumberOfMessagesVisible超过阈值。NumberOfMessagesReceived长时间为零。死信队列消息数大于零。用命令查看队列监控指标aws cloudwatch list-metrics --namespace AWS/SQS --dimensions NameQueueName,Valueorder-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 无服务器消费模式。分布式事务补偿方案与消息最终一致性。如果你在配置或代码运行中遇到其他问题欢迎在评论区留言也可以收藏本文备用。动手把示例代码跑一遍比囫囵吞枣看十篇文章更有用。