
微服务架构里服务之间的通信方式基本决定了整个系统的扩展边界和故障半径。AWS SNS 和 SQS 是 AWS 上最常用的一组消息服务SNS 负责把一条消息广播给多个订阅方SQS 负责把消息可靠地暂存起来让消费者按自己的节奏去拉取。这篇文章是系列的第一篇重点解决从零搭建、理解模型到跑通第一个发布订阅流程的问题适合刚开始接触 AWS 消息服务、或者已经在写微服务但还没把异步通信链路理顺的同学。其实很多人第一次接触 SNS 和 SQS 的时候最困惑的不是怎么创建资源而是这两个东西到底有什么区别。控制台里都能创建都能发消息看起来很像。但真正设计系统的时候选错模型会让后续的扩展、重试、消息顺序都变得别扭。这篇文章会先把模型讲清楚再带你把环境、权限、队列、主题、代码、验证全流程跑一遍最后给出一份我平时排查消息问题时一定会看的检查清单。读完这一篇你应该能独立搭出一个“生产者发布到 SNSSNS 扇出到 SQS消费者从 SQS 拉取处理”的最小链路。1. 微服务通信为什么绕不开消息队列1.1 同步调用在什么时候会变成瓶颈微服务之间最直接的通信方式是 HTTP 同步调用。服务 A 调用服务 B等 B 返回结果再继续往下走。这种模式写起来简单调试也直观但一旦系统规模上来问题就出现了。第一个问题是调用链被拉长。用户请求进来之后如果订单服务需要调用库存服务、支付服务、通知服务每个调用都占用一个线程或协程。任何一个下游服务变慢整个请求的响应时间都会被拖住。更麻烦的是下游服务一旦不可用上游服务的线程池会被大量占满最后出现连锁故障。第二个问题是高峰期流量难以削峰。比如电商秒杀场景一瞬间产生的订单量可能是平时的几十倍。如果所有请求都直接打到数据库或下游服务系统很容易被打垮。同步调用没有办法把峰值流量先存起来只能靠扩容硬抗而扩容又不可能按瞬时峰值一直预留。第三个问题是业务耦合。同步调用意味着服务 A 必须知道服务 B 的地址、接口参数、返回结构。服务 B 换了协议或者拆分成多个服务A 就要跟着改。微服务的数量越多这种强依赖关系就越难维护。消息队列解决的就是这三个问题异步化、削峰填谷、解耦。AWS 上的 SNS 和 SQS 就是这套思路的托管实现。1.2 SNS 和 SQS 在架构里各自承担什么角色SNS 全称 Simple Notification Service是发布订阅服务。一条消息发布到 SNS 主题之后会被推送给所有订阅者。订阅者可以是 HTTP 接口、Lambda 函数、邮件、手机号也可以是 SQS 队列。SQS 全称 Simple Queue Service是消息队列服务。消息被发送到队列之后消费者通过主动拉取的方式获取消息。每条消息在默认语义下只会被一个消费者拿到。用一句话区分SNS 解决“一条消息怎么发给多个接收方”SQS 解决“一条消息怎么暂存并交付给一个处理方”。在微服务架构里SNS 和 SQS 最常见的组合方式是“扇出模式”业务服务把消息发布到 SNS 主题主题下面挂多个 SQS 队列每个队列对应一个业务消费者。比如订单创建成功后订单服务发布一条 OrderCreated 消息库存服务、积分服务、通知服务各自监听自己的队列互不干扰。这种设计的好处是新增消费者不需要改动发布方的代码。以后想加一个数据分析消费者只需要在 SNS 主题下新增一个订阅把消息投递到新队列即可。2. 先把 SNS 和 SQS 的数据模型分清2.1 SNS 像广播SQS 像邮箱我习惯用一个生活化的类比来理解这两个服务。SNS 像广播电台。电台发出一条节目信号所有打开收音机的人都能收到。发送方不知道具体有多少人在听也不需要关心谁没有听到。在 AWS 里这个“收音机”就是订阅关系。SQS 像邮箱。你把一封信投进邮箱这封信会一直待在邮箱里直到收件人把它取走。同一个邮箱里的每一封信理论上只属于一个收件人。如果收件人暂时没空信不会丢会继续留在里面。这两种模型对应了完全不同的一致性语义。SNS 的投递语义是“尽力而为”。消息发布到主题后SNS 会尝试推送给每个订阅者但如果有订阅者不可达消息可能会丢失。当然如果你订阅的是 SQS 队列消息会先进入队列再由队列保证持久化这时候可靠性就由 SQS 来兜底。SQS 的语义是“至少一次投递”。消息进入队列后会被持久化存储消费者拉取后需要显式删除消息。如果消费者处理失败没有删除消息会重新可见再次被消费。这个机制意味着同一个业务消息可能被处理多次所以消费端的幂等设计必不可少。2.2 三种最常见的组合关系刚上手的人往往纠结“到底用 SNS 还是 SQS”其实很多场景里两者是一起用的。第一种组合只用 SQS。适合“一个生产者、一个消费者”的场景。比如后台任务处理应用把任务消息塞进队列Worker 进程拉取执行。这时候不需要广播SQS 自己就够。第二种组合只用 SNS。适合“一个生产者、多个订阅方”且订阅方是 Lambda 或 HTTP 端点的场景。比如系统告警通知一个监控服务把告警消息发布到主题邮件订阅者收到邮件短信订阅者收到短信Lambda 订阅者执行自动化处理。第三种组合SNS 加多个 SQS。这是微服务架构里最推荐的方式。生产者只和 SNS 主题打交道不感知下游有哪些消费者。每个消费者拥有自己的队列消费者宕机不影响其他消费者消息会留在队列里等恢复。我一般建议项目初期的架构设计阶段就把消息模型画出来谁产生消息谁消费消息一个消息需要给几个业务方每个业务方是否可以独立延迟处理。画完这些选型自然就清楚了。3. 环境准备账号、CLI、IAM 权限3.1 本地环境需要准备什么跑通本文的示例需要准备四样东西第一一个 AWS 账号。如果只是学习建议使用新账号或者单独的区域避免影响已有业务资源。云服务都是按量计费SNS 和 SQS 有免费额度但测试创建的资源记得最后清理。第二本地的 AWS CLI。CLI 不是必须的控制台也能完成所有操作但后续如果你要写自动化脚本CLI 和 SDK 会更顺手。CLI 安装后运行aws --version确认版本正常。第三Python 3 环境和 boto3 库。本文的代码示例使用 Python 的 boto3 SDK这也是 AWS 生态里最常用的开发方式。安装命令很简单pip install boto3第四一组访问密钥。密钥用于本地程序调用 AWS API。创建密钥时建议只赋予最小权限不要直接给 AdministratorAccess。注意密钥文件下载后只出现一次妥善保存。不要把 Access Key 提交到代码仓库也不要在博客、聊天工具里明文展示。3.2 IAM 权限最小化配置我见过很多新手图省事直接在本地用根用户的密钥跑程序。短期测试没问题但一旦泄露影响的是整个账号。更稳妥的做法是创建一个专门的 IAM 用户只授予本次实验需要的权限。如果需要创建和管理 SNS、SQS 资源可以先用一个偏宽松的策略下面只是示例落地时要按自己账号的情况裁剪{ Version: 2012-10-17, Statement: [ { Effect: Allow, Action: [ sns:CreateTopic, sns:Subscribe, sns:Publish, sqs:CreateQueue, sqs:SendMessage, sqs:ReceiveMessage, sqs:DeleteMessage, sqs:GetQueueAttributes ], Resource: * } ] }实际生产环境里权限应该进一步收紧到指定资源 ARN。比如 Program A 只能往某个主题发布只能消费某个队列。AWS 的 IAM 策略支持按资源 ARN 限制这是生产环境必须做的。3.3 区域选择和公共网络前提SNS 和 SQS 都是区域级服务。你在哪个区域创建主题和队列消息就存储在哪个区域。区域之间是完全隔离的不同区域的客户端不能直接访问另一个区域的资源。选择区域时优先考虑两点一是你的业务资源所在区域二是目标用户所在区域。比如业务部署在 ap-southeast-1消息服务也建在同一个区域跨区域调用会增加延迟和成本。本地程序调用 AWS API 需要能正常访问 AWS 的公共端点。如果你在云服务器上运行示例需要确认安全组和系统防火墙允许出站 HTTPS 访问 443 端口。AWS 的 API 全部走 HTTPS所以只需要确保 443 出站是通的。注意如果本地网络访问 AWS 控制台或 API 不稳定先检查网络连通性和代理设置而不是贸然改代码重试。网络问题表现成超时报错时最容易让人误判成代码问题。4. 创建第一个 SQS 队列和 SNS 主题4.1 控制台创建 SQS 队列登录 AWS 控制台搜索 SQS进入队列页面点击“创建队列”。这里先选择标准队列不要选 FIFO。标准队列吞吐高但消息顺序不保证FIFO 队列保证先进先出但每秒请求次数有限制。学习阶段两种队列差异主要体现在顺序和吞吐上后文会单独讲。创建队列时需要关心几个参数队列名称建议带上环境标识比如dev-order-queue。可见性超时默认 30 秒。消费者拉取一条消息后消息会进入“不可见”状态。如果 30 秒内消费者没有删除消息消息会重新变成可见状态可能被其他消费者再次拉到。这个参数要根据消费者处理一条消息的时长来设置。消息保留期默认 4 天。消息在队列里最多保留多久超过就自动删除。如果业务允许延迟处理可以调大如果对时效性要求高可以调小。最大消息大小默认 1 KB实际可以调到 256 KB。超过大小的消息会发送失败。创建完成后控制台会显示队列的 URL 和 ARN。URL 是客户端访问队列用的ARN 是授权和订阅用的两者不要混淆。4.2 控制台创建 SNS 主题搜索 SNS进入主题页面点击“创建主题”。类型选择“标准”输入主题名称比如dev-order-topic。创建完成后记录主题 ARN形如arn:aws:sns:ap-southeast-1:123456789012:dev-order-topic接下来把刚创建的 SQS 队列订阅到这个主题。在主题详情页点击“创建订阅”协议选择“Amazon SQS”端点输入队列 ARN。提交后订阅状态会变成已确认因为 AWS 会自动完成 SQS 订阅确认。订阅完成后可以做一个最简单的验证在主题详情页点击“发布消息”输入消息内容并发布然后回到 SQS 队列页面查看消息数。如果队列的消息可用数变成 1说明链路已经通了。4.3 用 CLI 创建并验证如果你打算写自动化脚本用 CLI 创建会更高效。下面是三条常用命令。创建队列aws sqs create-queue \ --queue-name dev-order-queue \ --region ap-southeast-1 \ --attributes VisibilityTimeout30,MessageRetentionPeriod345600创建主题aws sns create-topic \ --name dev-order-topic \ --region ap-southeast-1订阅队列到主题aws sns subscribe \ --topic-arn arn:aws:sns:ap-southeast-1:123456789012:dev-order-topic \ --protocol sqs \ --notification-endpoint arn:aws:sqs:ap-southeast-1:123456789012:dev-order-queue每条命令执行成功后都会有返回结果。创建队列返回 QueueUrl创建主题返回 TopicArn订阅返回 SubscriptionArn。如果某一步没有返回优先看报错信息里的 AccessDenied 或 InvalidParameter。5. 用 boto3 跑通发布、订阅、消费全流程5.1 初始化客户端代码里使用 boto3 客户端时需要指定区域和凭证。推荐的方式是使用环境变量或 AWS 配置文件而不是在代码里硬编码密钥。import boto3 sqs boto3.client(sqs, region_nameap-southeast-1) sns boto3.client(sns, region_nameap-southeast-1)如果本地配置了aws configure上面这段代码可以不传密钥直接运行。boto3 会按顺序读取环境变量、本地配置文件、IAM 角色等凭证来源。5.2 发布消息到 SNS发布消息的核心方法是sns.publish。最简单的调用只需要传入 TopicArn 和 Message。response sns.publish( TopicArnarn:aws:sns:ap-southeast-1:123456789012:dev-order-topic, MessageHello from microservice A, ) print(response[MessageId])返回结果里会有一个 MessageId可以把它当成这条消息的唯一标识。后续排查时如果需要确认某条消息是否发送成功这个 ID 是关键线索。实际业务里消息内容通常是 JSON 字符串。建议在 Message 里传结构化的 JSON而不是纯文本。例如import json message_body { event_type: order.created, order_id: 202606010001, user_id: u_10086, amount: 199.9, } sns.publish( TopicArntopic_arn, Messagejson.dumps(message_body, ensure_asciiFalse), )这里的关键点是ensure_asciiFalse避免中文被转义成\uXXXX影响后续日志可读性。SNS 发布时还可以设置 MessageAttributes用来给消息附加元数据。订阅过滤策略就是基于消息属性来判断的后文会讲到。5.3 从 SQS 拉取消息并删除消费者从 SQS 拉取消息使用sqs.receive_message。response sqs.receive_message( QueueUrlhttps://sqs.ap-southeast-1.amazonaws.com/123456789012/dev-order-queue, MaxNumberOfMessages10, WaitTimeSeconds20, )MaxNumberOfMessages 控制一次最多拉多少条范围是 1 到 10。WaitTimeSeconds 是长轮询时间设为大于 0 的值可以避免空轮询造成的大量 API 请求。长轮询的意思是如果队列里没有消息请求会保持连接等待直到有消息或者达到超时时间。拿到消息后立刻打印并处理处理成功后再删除。for message in response.get(Messages, []): body message[Body] receipt_handle message[ReceiptHandle] print(fReceived: {body}) # 业务处理逻辑 process_message(body) sqs.delete_message( QueueUrlqueue_url, ReceiptHandlereceipt_handle, )ReceiptHandle 是这条消息本次接收的唯一凭证。删除消息时必须传它不传 MessageId。这个细节很容易踩坑新手经常拿着 MessageId 去删消息结果一直报错。删除消息完成整个“发布到 SNS、扇出到 SQS、从 SQS 消费”的链路就闭环了。5.4 完整的最小示例把上面的步骤串成一个可运行的脚本建议按这个结构组织import json import boto3 REGION ap-southeast-1 TOPIC_ARN arn:aws:sns:ap-southeast-1:123456789012:dev-order-topic QUEUE_URL https://sqs.ap-southeast-1.amazonaws.com/123456789012/dev-order-queue sns boto3.client(sns, region_nameREGION) sqs boto3.client(sqs, region_nameREGION) def publish_order_event(order_id: str, user_id: str): body { event_type: order.created, order_id: order_id, user_id: user_id, } resp sns.publish( TopicArnTOPIC_ARN, Messagejson.dumps(body, ensure_asciiFalse), ) return resp[MessageId] def consume_messages(max_num: int 5): resp sqs.receive_message( QueueUrlQUEUE_URL, MaxNumberOfMessagesmax_num, WaitTimeSeconds20, ) for msg in resp.get(Messages, []): print(f处理消息: {msg[Body]}) sqs.delete_message( QueueUrlQUEUE_URL, ReceiptHandlemsg[ReceiptHandle], ) if __name__ __main__: message_id publish_order_event(202606010001, u_10086) print(f已发布消息, MessageId{message_id}) consume_messages()运行后先看两件事发布时是否返回 MessageId消费时是否能打印出消息内容。两步都成功说明基础链路正常。如果消费端始终拿不到消息先回到控制台看队列里有没有消息积压有积压说明发布和订阅没问题问题出在消费端没有积压说明发布或订阅环节就已经断了。6. 事件驱动架构里最常见的组合SNS 扇出到多个 SQS6.1 扇出模式的业务价值单队列单消费者的模式只能解决“一个任务一个处理方”的问题。微服务场景里一个业务事件往往需要多个服务同时响应。以订单创建为例订单服务只需要负责订单本身的状态流转但订单创建成功后优惠券服务要发券积分服务要加积分通知服务要发站内信。如果订单服务直接调用这三个服务就回到了同步耦合的老路。扇出模式的做法是订单服务发布一条消息到 SNS 主题主题下面挂三个 SQS 队列分别对应优惠券、积分、通知三个消费者。每个消费者只关心自己的队列。这种模式下新增消费者不需要改动订单服务代码。比如后续要增加一个“订单数据进数仓”的消费者只需要新建一个队列并订阅到原主题订单服务完全无感。6.2 订阅过滤策略怎么用扇出模式有一个衍生需求有些消费者只关心特定类型的消息。比如订单状态有 order.created、order.paid、order.cancelled通知服务想全部处理但积分服务只关心 order.paid。如果每个消费者都是全量消费消费者拿到消息后自己再判断类型代码里就多了一层 if-else。更好的做法是使用 SNS 的订阅过滤策略让消息在 SNS 层面就被分流。使用过滤策略时发布方必须给消息设置消息属性。以 Python 为例sns.publish( TopicArntopic_arn, Messagejson.dumps(body, ensure_asciiFalse), MessageAttributes{ event_type: { DataType: String, StringValue: order.paid, } }, )订阅端创建订阅后可以给订阅关联一个 JSON 格式的过滤策略{ event_type: [order.paid] }设置了过滤策略的订阅只会收到 event_type 等于 order.paid 的消息。其他消息不会进入对应的 SQS 队列。过滤策略支持的匹配方式包括精确匹配、前缀匹配、数值范围等。一开始不需要把策略写得很复杂先掌握精确匹配就够用了。6.3 消息字段和日志上下文设计消息体时建议包含几个公共字段事件 ID、事件类型、事件时间、业务主键、其他业务数据。{ event_id: evt_0001, event_type: order.paid, occurred_at: 2026-06-01T10:30:00Z, business_key: order_202606010001, data: {} }event_id 用于全链路追踪。消费者日志里记录 event_id 后如果消息处理异常可以拿这个 ID 回查发布方日志。business_key 用于幂等判断消费者可以用它做去重。7. 消息丢失、重复和堆积怎么排查7.1 先看现象再分层定位消息系统出问题时最常见的现象是三类消息没了、消息重复了、消息堆积了。消息没了先确认是发布侧没发出去还是消费侧没收到。发布侧看 SNS 的 Publish 调用是否返回 MessageId返回了说明 SNS 已经接收。再看 SQS 控制台队列的“消息可用数”是否增长。如果增长但消费者没处理问题在消费端如果没增长问题在订阅关系或过滤策略。消息重复了先确认消费者代码里删除消息是否一定执行。只要消费者拉取消息后网络中断、进程重启、删除超时消息就会重新可见。重复不是 bug是至少一次投递语义的必然结果。生产环境要做的是让处理逻辑幂等。消息堆积了先看堆积发生在哪个队列。单个消费者处理速度跟不上要么扩消费者实例要么调大每次拉取数量。如果消费端报错不断消息反复被拉取并超时堆积会越来越严重这时候优先修消费端而不是盲目扩容。7.2 DLQ 死信队列是必须提前配的消息处理失败时默认行为是重新投递。如果消费者代码有 bug消息会进入“拉取-失败-超时-重新可见”的循环。这种消息会持续消耗资源而且会阻塞队列里后面的消息。标准做法是配置死信队列。当一个消息被消费失败的次数超过阈值SQS 会把它投递到关联的 DLQ。消费者代码只需要处理主队列处理失败的消息会被自动转移。创建死信队列时需要给源队列配置 RedrivePolicy。下面是 CLI 示例aws sqs set-queue-attributes \ --queue-url https://sqs.ap-southeast-1.amazonaws.com/123456789012/dev-order-queue \ --attributes { RedrivePolicy: {\deadLetterTargetArn\:\arn:aws:sqs:ap-southeast-1:123456789012:dev-order-queue-dlq\,\maxReceiveCount\:5} }maxReceiveCount 设为 5 的意思是一条消息被拉取 5 次仍未删除就转到 DLQ。这个值不要设太小因为有些消息处理失败可能是临时性的比如数据库闪断也不要设太大否则问题消息会长时间占着队列。7.3 可见性超时和重试顺序消息处理超时的根源通常是可见性超时设置不合理。消费者拉取消息后如果处理时间超过 VisibilityTimeout消息会提前变回可见状态被其他消费者拉到导致同一条消息被并发处理。排查时先看消费端平均处理一条消息要多久然后把 VisibilityTimeout 设成平均耗时的 3 到 5 倍。比如平均 5 秒处理完可见性超时设为 30 秒比较安全。另外要注意SQS 标准队列不保证顺序。如果业务要求消息严格按顺序处理比如同一个订单的状态不能乱序只能改用 FIFO 队列。FIFO 队列按消息组 ID 保证组内顺序但吞吐量受限需要单独做性能评估。7.4 推送类消费者的失败处理如果你用 SNS 直接触发 Lambda而不是经由 SQS失败语义又不一样。SNS 推送到 Lambda 失败时会自动重试几次但还是会存在最终丢弃的可能。如果业务要求消息不丢最稳妥的方案仍然是 SNS 加 SQS 加 Lambda 或加 WorkerSNS 保证扇出SQS 保证持久化消费者保证处理进度。消息先落队列再进业务逻辑可靠性会高很多。8. 生产环境落地的经验清单8.1 命名、标签和资源治理消息服务的资源一旦多起来命名混乱是最先暴露的问题。建议统一用“环境-业务域-用途”的格式比如prod-order-queue、staging-payment-topic。所有资源建议打上标签比如 Environment、Team、CostCenter。AWS 的标签系统能辅助做成本拆分和资源归属后续对账时省很多事。8.2 消费幂等是硬性要求SQS 提供的是至少一次投递语义消费者必须假设同一条消息可能收到多次。幂等方案有以下几种利用数据库唯一键去重。比如消息里的 business_key 在表里建唯一索引重复插入会直接失败程序捕获冲突后视为成功。利用 Redis 做短时间去重。处理前先 SETNX设置一个短过期时间只有第一次能设置成功。利用业务状态机判断。订单已经从未支付变成已支付再收到一条支付消息时直接跳过。幂等不能靠“概率上应该不会重复”来掩盖必须从代码层面兜住。8.3 批量处理和并发控制单条消息处理吞吐通常不够。SQS 支持批量接口SNS 也支持批量发布。批量发送消息到 SQS 使用send_message_batch一次最多 10 条。批量消费时receive_message的 MaxNumberOfMessages 设置为 10处理完一批后统一删除。并发控制要分两层看。消费者进程内部可以用线程池或协程并发处理多条消息但并发数不要一开始就拉满。先观察目标数据库或下游接口的承受能力再逐步调大。如果消费者是以 EC2 或 ECS 任务运行横向扩容时要注意队列里的消息会被多个实例同时拉取重复处理的可能性会上升。幂等没做好前不建议开太多实例。8.4 成本和监控SNS 和 SQS 的计费主要是按请求次数和数据传输量。避免成本失控的方法不要在代码里做短轮询。WaitTimeSeconds 必须大于 0否则每条空消息都会产生请求费用。合理设置消息保留期。消息积压太久会导致存储费用增加而且积压消息往往是消费端异常的信号。配置 CloudWatch 告警。重点盯四个指标队列深度、旧消息时间、DLQ 消息数、发布失败数。队列深度持续上升或 DLQ 出现消息应该立刻查看消费端日志。8.5 一套推荐的上手顺序最后整理一下我自己上手这套架构时采用的顺序也建议你按这个节奏走先用控制台创建主题和队列手动发布一条消息验证订阅链路。然后写一个最简单的 Python 脚本发布和消费各一条消息确认 SDK 调用方式正确。再扩展成扇出模式创建第二个队列验证不同消费者各收各的消息。之后加上过滤策略、DLQ、批量接口补齐生产需要的能力。最后配置监控告警和日志把资源清理规则定下来。这套顺序的好处是每一层都能验证出了问题很容易定位。不要一上来就写几百行代码把所有功能都堆在一个脚本里。消息系统出问题时最难的不是代码逻辑而是你不知道消息卡在哪一层。这一篇讲的是从零搭建最小链路和核心概念。下一部分可以继续聊更进阶的主题FIFO 队列的顺序保证、SNS 推送的限流和重试策略、消费者容器化部署时的优雅停机以及消息系统在微服务架构里的更多落地细节。如果你刚接触 AWS 消息服务建议先把这一篇里的示例跑通再往后看。