用SNS与SQS重塑微服务通信:原理、链路搭建与生产避坑 大概三年前我接手过一个微服务改造的评估任务。系统里几个服务之间用的是最直接的 HTTP 调用支付完成后要调库存服务扣库存再调通知服务发消息。平时流量不高一切正常。第一次大促压测下游一个接口慢了 200 毫秒结果上游所有线程都堵在等待里数据库连接先被耗尽紧接着整条调用链跟着雪崩。那是我第一次意识到微服务之间的通信方式比服务内部的代码设计更能决定系统的生死。后来几个项目里陆续用上 AWS 的 SNS 和 SQS回头看这套方案真正有价值的不是“能发消息”这个功能而是它改变了微服务之间协作的结构把硬邦邦的同步调用变成了可以缓冲、可以广播、可以重试的异步事件流。这篇文章是系列的第一篇。我不会把 SNS 和 SQS 的 API 参数全部列一遍而是先把这两个服务在微服务架构里的定位讲清楚然后给出一条可以照着跑通的最小链路再讲几个生产环境里最容易踩的坑。看完之后你应该能判断自己的系统该用队列还是该用事件总线或者两个一起用。1. 微服务之间为什么需要“中间人”先理解耦合与故障扩散很多人把 SNS/SQS 当成“消息中间件”这个说法没有错但它忽略了一个更重要的点这两个服务真正解决的是微服务之间的结构性耦合问题。耦合不只是在代码层面也体现在服务之间的交互方式上。1.1 同步调用的三个隐性问题先看一个最常见的场景订单服务创建订单后需要通知库存服务扣减库存再调用会员服务增加积分。用 HTTP 同步调用时表面上是两个独立的服务实际上订单服务要被迫等待库存和会员两个服务都返回结果才能给用户一个成功响应。这里藏着三个在低流量时很难发现的问题故障扩散。库存服务如果出现数据库慢查询响应时间从 100ms 变成 5 秒订单服务的线程池很快就会被打满。订单服务又撑不住调用它的上游网关也开始积压请求。一个服务的问题像多米诺骨牌一样向下传递。流量无缓冲。系统平时能抗住每秒 200 个请求大促时流量变成每秒 2000 个。如果所有请求都直接打到下游服务下游只能被动接受要么扩容要么限流。扩容需要时间限流则直接损失业务。发布耦合。订单服务每次发版只要涉及调用关系都要和下游团队同时测试、同时发布。微服务本来是为了独立部署结果同步调用又把大家绑在了一起。这三个问题本质上都源于同一个设计假设假设下游一定可用而且一定在预期时间内返回。在复杂的分布式环境下这个假设每隔一段时间就会被打破一次。1.2 引入中间人之后发生了什么SNS 和 SQS 做的事情是在调用链里插入一个“中间人”。订单服务不再直接调用库存服务而是把“订单已创建”这个事实写进消息发到队列或主题。之后谁来消费、什么时候消费、消费完之后做什么都变成下游自己的事。这个变化带来的实际收益是问题同步调用异步消息SNS/SQS下游故障上游被迫等待线程被耗尽消息积压在队列里上游直接返回成功突发流量下游直接被冲垮队列充当缓冲池消费速率决定处理速度服务发布上下游必须协调联调消息契约稳定消费方可以独立升级扩展能力需要同时扩展整条链路上的服务只需扩展消息消费者的数量这里要强调一个关键认知引入 SNS/SQS 并不是让微服务变快而是让系统在故障和峰值面前不会立刻死掉。消息积压是相对可控的状态接口超时导致线程池耗尽则是崩溃的前兆。两者有本质区别。2. SNS 和 SQS先分清广播与队列别混为一谈不管在 AWS 官方文档里还是社区讨论里SNS 和 SQS 经常被放在一起说很多文章甚至用“SNS/SQS”这个合体词指代一整类方案。但在实际架构设计时这两个服务的语义完全不同用错地方后续会非常别扭。2.1 SQS点对点的任务队列SQS 的核心模型是一个队列。生产方把消息放进队列消费方从队列里拉取消息。消息被消费后消费方要显式删除否则这条消息会再次出现在队列里。这个模型适合什么场景适合“把一件需要被处理的任务交给一个处理者”。例如订单服务往队列里放一条“扣减库存”的消息任何消费者取到这条消息后只需要完成扣减动作。任务本身不会重复执行第二次队列也不关心这个消息要不要广播给其他服务。SQS 的关键特征是消费方主动拉取。队列不负责把消息推给消费者消费者按照自己的节奏来取。这意味着即使消费者短暂不可用消息还在队列里积压不会丢失也不会把消费者压垮。2.2 SNS一对多的事件广播SNS 的核心模型则完全不同。它不是一个队列而是一个主题Topic。生产者把事件发布到主题上主题会把这条事件复制并推送给所有订阅了该主题的终端。终端可以是 SQS 队列、Lambda 函数、HTTP/HTTPS 端点、电子邮件等。一个订单创建事件发布到 SNS 主题后负责积分、客服、数据分析的服务都能各自订阅收到同一份事件各自处理。这就像广播电台一个主播说话所有听众都能听到但每个听众如何反应主播并不关心。SNS 适合的是“把事实告诉所有人”的场景而不是“把任务分给一个处理者”。2.3 两者组合扇出Fan-out才是微服务里的常见形态真正让微服务团队兴奋的不是单独使用 SQS 或 SNS而是把两者组合起来SNS 主题订阅多个 SQS 队列每个队列对应一个独立消费方。一个订单创建事件发到 SNS 主题后主题自动把事件推送给库存队列、会员队列、数据分析队列。每个下游服务只消费自己的队列互不干扰。一个服务处理失败只影响自己的队列不会阻塞其他服务。这就是典型的扇出架构。可以这样记忆SQS 是请了一个搬运工把任务从仓库送到指定的人手里。SNS 是一个大喇叭喊一嗓子让所有感兴趣的人都听到。SNS SQS 是大喇叭前面挂了多个信箱每个团队从自己的信箱里取信各干各的。3. 把第一条消息跑通最大可用链路理解概念之后先别急着设计复杂拓扑。把一条消息从 SQS 发出去再用一个消费者收下来确认每一步的表现再慢慢增加复杂度。3.1 创建标准队列并在控制台验证登录 AWS 控制台进入 SQS 页面点击创建队列。队列类型选择“标准队列”先不要勾选“FIFO 队列”。FIFO 队列有顺序保证但吞吐上限低得多第一步先用标准队列跑通流程。创建完成后控制台会显示队列 URL。这个 URL 就是后续所有 API 调用的定位标识。在控制台消息页面可以直接发送一条测试消息。发送成功后队列的“可用消息数”会变成 1。这里最容易忽略的是理解“消息被接收”和“消息被删除”的区别。在 SQS 里消费者拉取消息后消息不会立刻消失而是进入一种“不可见”状态持续一段时间即可见性超时。如果消费者在这段时间内没有显式删除消息消息会重新变得可见被其他消费者再次拉到。这既是 SQS 保证消息最终被处理的方式也是重复消费的来源。3.2 用 Lambda 做消费者最小实现创建一个简单的 Lambda 函数并为它配置 SQS 触发源。这是目前最省事的消费方式因为消息的拉取、删除、重试都由 AWS 托管处理开发人员只需要写业务逻辑。下面是一个常见的 Python 处理函数结构import json def handler(event, context): # 每条 recode 对应一条 SQS 消息 for record in event[Records]: body json.loads(record[body]) print(f处理消息{body}) # 在这里写你的业务处理逻辑 # 例如更新数据库、调用下游接口、写日志等 return {statusCode: 200}关键点在于返回值。如果 Lambda 函数执行成功AWS 会自动删除已经处理完的 SQS 消息如果函数抛出异常或执行失败消息不会被删除会根据配置的重试策略再次进入函数。对于第一步验证直接把消息内容和日志打印出来即可。在创建 Lambda 触发器时需要注意一个配置批量大小Batch size。默认是 10也就是一次把 10 条消息给到函数。先不要调整这个值用默认值跑一遍重点观察 CloudWatch 日志里是否能打印出刚才发送的测试消息。3.3 从控制台到命令行建立可复用的发送方式控制台发送消息只适合验证界面。一旦要真正对接业务系统最常见的还是通过 SDK 发送。用 boto3 发送消息的代码很直接import boto3 sqs boto3.client(sqs, region_nameap-northeast-1) response sqs.send_message( QueueUrlhttps://sqs.ap-northeast-1.amazonaws.com/123456789012/my-queue, MessageBody{order_id: 10001, sku: A001} ) print(response[MessageId])这里最需要注意的是 QueueUrl 是当前环境专属的不同的区域、账号、队列名都会生成不同的 URL。实际编码时通常会把队列 URL 放到环境变量或配置中心而不是硬编码在代码里。先跑通单条消息再考虑批量。这不是保守而是为了让你能够清晰地区分问题发生在生产端、队列端还是消费端。4. 不调参等于裸奔几个决定成败的关键配置很多人第一次跑通 SQS 链路之后觉得这事很简单发消息收消息删消息。等到上了生产环境、开始处理真实业务才发现消息丢失、重复、积压、乱序每一个问题都和生产环境的具体参数配置有关。4.1 可见性超时消费者处理多久才算合理可见性超时Visibility Timeout是 SQS 最重要的参数之一。它决定了一条消息被消费者取走后在多长时间内对队列中的其他消费者不可见。如果这个值设置太短消费者还在处理中消息就重新变可见被另一个消费者拉到导致同一条业务消息被两个消费者同时处理。如果设置太长消费者崩溃后消息要等待很久才会重新变成可消费状态拖慢了整个系统的恢复速度。在 Lambda 触发 SQS 的场景里AWS 官方给过一个经验值可见性超时的时间最好设置成函数超时时间的 6 倍。例如函数超时是 30 秒可见性超时就设为 180 秒。这样即使 Lambda 内部发生重试也不会在函数还没结束前就把消息暴露给其他消费者。4.2 死信队列系统必须预留的兜底所谓死信队列DLQ是把反复处理失败的消息从主队列中转存到另一个队列。它的价值不是避免失败而是保护主队列不被坏消息堵死。举个例子。消费者从队列里取到一条消息消息体是一个 JSON但由于某种原因这个 JSON 里的某个字段始终解析失败。如果不设置 DLQ消费者会反复拿到这条消息反复失败。在一个活跃业务里坏消息的比例可能只有百万分之一但在高吞吐场景下百万分之一也足以让某几个消费者实例持续空转还会让指标看起来像是在“重试”。配置上主要是两个值Maximum receives 之后进入 DLQ比如设为 3表示一条消息最多被接收 3 次之后转入死信队列。DLQ 本身也是一个 SQS 队列。我见过不少团队上了生产环境才发现忘记配置 DLQ消息反复失败却无从查起。第一版就带上 DLQ成本很低收益很高。4.3 重复消费、乱序与 FIFO 队列标准 SQS 队列提供的是至少一次交付at-least-once不保证不会重复。原因是消息的删除和写入之间存在一个时间窗消费者处理成功后删除请求如果因为网络原因丢失消息在可见性超时后会再次被消费。这一点必须在设计数据写入时提前想好。消费者如果是对数据库做操作尽量让业务处理逻辑具备幂等性或者在处理前检查去重表、记录消息 ID。不要在线上排查时才后悔没有设计幂等。对于强顺序要求的场景标准队列并不合适。SQS 的 FIFO 队列提供了严格的消息顺序和仅一次处理within a message group但代价是吞吐上限远低于标准队列而且发送时需要指定消息组 ID。FIFO 适合银行对账、订单状态流转这类必须保持有序的场景日志推送、缓存刷新、事件广播这类场景用标准队列就够了。4.4 长轮询小参数隐藏的性价比SQS 消费有两种轮询方式。短轮询是消费者发起请求后如果队列里没有消息立刻返回空结果长轮询是请求到了队列如果没有消息连接保持一段时间等待消息到达默认最多 20 秒。长期跑生产环境的团队一定要把 ReceiveMessage 的等待时间设成长轮询。它最大的收益不是减少响应时间而是显著减少空轮询请求的数量从而降低 SQS 的费用和无关的 API 调用。这个参数在实际项目里经常被忽略因为短轮询在小流量下也能工作但服务器多了以后每个消费者每秒轮询一次积少成多就是一个不小的成本。5. 放进微服务架构里三种常见拓扑SNS 和 SQS 的价值最终要放进真实架构里才能体现。这里的三种拓扑按复杂度从低到高排列也是很多团队逐步演进的过程。5.1 点对点订单服务到履约服务最基础的做法订单服务创建订单后把“待履约订单”的消息发到一个 SQS 队列履约服务作为消费者拉取消息并处理。这个拓扑解决的问题是削峰填谷。大促时订单量突然增高订单服务只需要把消息快速写入队列并返回履约服务按照自己的最大消费能力处理不会因为流量波动被打垮。即使履约服务短暂重启消息还在队列里恢复后继续消费。这适合业务中只有一个下游需要对某个事件做出反应并且不需要广播的场景。系统复杂度最低也最容易排查。5.2 扇出一个事件多个服务各自响应当同一个事件需要被多个服务消费时直接调多个接口会再次回到“同步调用导致耦合”的问题。SNS 多个 SQS 队列的结构让每一条下游链路都是独立的。例如订单创建事件SNS 主题order-events订阅者 A库存队列 → 库存服务订阅者 B积分队列 → 会员服务订阅者 C分析队列 → 数据管道库存服务处理慢不影响积分服务积分服务消费失败也不影响库存扣减。每个团队只管自己的队列和消费逻辑不需要知道还有哪些服务订阅了同一个事件。这种模式在很多场景里都会出现。例如物联网场景里的固件升级任务后端把升级指令发布到主题不同区域的升级服务各自订阅并推进 OTA 任务又比如 AI 推理任务完成后把结果事件广播给多个下游做缓存刷新、通知和报表统计。理解扇出模式之后很多看起来复杂的协作需求本质上都是同一个“一个事件多个消费者”的问题。5.3 按业务领域拆分消费边界到了这一步系统已经不是“一个事件对应几个服务”了而是在架构层面把事件流作为服务之间的主通信协议。团队可以按限界上下文来划分消息的消费边界。订单领域发布订单事件库存领域消费并更新库存支付领域消费并生成账单。任何两个服务之间不直接产生同步依赖而是各自订阅上游领域发布的事件。这本质上是事件驱动架构在微服务里的实现。这种拓扑最灵活但要求团队有很强的消息契约管理能力。事件结构的变更需要版本管理消费方需要容忍消息字段的兼容性演进。如果团队没有形成这样的协作规范贸然把所有通信都改成事件反而会乱成一团。6. 消息不见了、重复了排查链路怎么走接手别人的系统时最常遇到的排查场景就是两类消息丢了或者消息重复处理。下面是一条经过多次验证的排查链路。6.1 从现象到分层排查先看现象再看输入再看消费端。步骤一般是确认消息是否进了队列。在控制台看队列的“可用消息数”和“进入/出站消息数”。如果发送 API 返回成功但消息数一直是 0重点检查发送时指定的队列 URL 是否和生产环境的队列一致。检查消费者是否正常拉取。如果消息数一直在涨说明消费者没有消费或消费者的权限配置有问题。进入 CloudWatch 查看 Lambda 的执行日志确认是函数没有触发还是函数触发了但一直抛异常。检查可见性超时和重试次数。消息被反复消费但看不到业务失败记录可能是可见性超时太短函数还在处理中消息就已经重新可见。检查是否进入了 DLQ。如果主队列消息数下降但 DLQ 里出现了消息说明业务处理连续失败。这时候优先看 DLQ 里消息体本身往往问题一目了然。检查消费端代码的删除逻辑。使用 Lambda 时函数成功自动删除如果自己写消费者需要确认是否在业务处理成功之后才调用 DeleteMessage。一个常见的错误是收到消息后立刻删除再执行业务逻辑结果业务失败消息已经没了。6.2 重复消息怎么查重复消息在标准队列里是预期内行为。看到重复不要急着怀疑 AWS 丢消息或有 bug而是优先检查业务是否做到了幂等。排查重复时有三个方向消费者是否在删除前崩溃。如果是重试必然导致重复。生产者是否重复发送。网络超时后客户端重试生产者可能会把同一条业务消息发两次。消息体里是否有全局唯一 ID。如果消息带有业务编号可以在消费端做去重如果消息体里没有合适的主键就要考虑在发送前为消息生成一个唯一 ID 放入消息属性消费者用这个 ID 做幂等判断。6.3 日志和追踪能力在生产环境里单靠控制台看消息数是不够的。我建议每一条消息的发送、接收、处理完成、进入 DLQ都要打印日志并带上消息 ID。涉及跨服务链路时把业务 ID 和消息 ID 一起记录排查时按业务 ID 串起整条链路。如果没有追踪能力出现积压时你只能看到一个队列数字在涨很难定位是哪个消费者处理变慢了。这一步算不上高深但很多团队恰恰倒在“没打日志”上。7. 别把 SNS/SQS 当万能药适用边界写到这里应该强调一个反向的问题SNS/SQS 不是所有微服务通信问题的答案。它在解决一类问题的同时也引入了一些新的约束。7.1 不适合低延迟、强一致性的请求场景如果你的业务需要用户发起请求后在几百毫秒内得到确定性的处理结果比如查询订单状态、支付确认页的即时校验异步消息并不合适。消息积压、消费者调度、重试机制都会让响应时间变得不可控。这类场景应该保留同步接口或者使用 API Gateway Lambda 这类同步服务。SNS/SQS 的本质是“接受延迟和不确定性”。它把“立即出结果”的承诺换成了“最终会处理”的保证。数据结构、用户体验、业务流程都必须适配这个心智模型。7.2 不适合大数据包和强顺序高频场景SQS 单条消息体有大小上限默认是 256KB。如果业务要在消息体里直接塞图片、视频或几百 KB 的日志就不太合适。常规做法是消息体里只放文件路径或对象 key真正的大数据存到 S3消费者再去拉取。强顺序场景也必须谨慎。标准队列本身就是无序、可重复的FIFO 队列虽然能保证顺序但吞吐上限低。如果业务需要海量有序事件可能需要考虑 Kinesis Data Streams 这类流式方案。7.3 只适合接受“最终一致”的链路使用 SNS/SQS 的服务间关系在数据层面永远是最终一致的。订单服务和库存服务之间不会立刻一致中间隔着队列和消息处理时间。需要强事务保证的环节不能靠消息来硬撑必须在业务流程或补偿机制里设计清楚。落在最后的一句话从第一次压测雪崩到后来在项目里稳定使用 SNS/SQS我最大的感受是微服务架构的难点从来不是拆服务而是服务之间如何协作。SNS/SQS 提供了一种成本可控的协作方式让服务之间不再彼此拖累让每一个团队都能按照自己的节奏演进。但这套方案的心智模型需要一阵子适应它要求你把“调用成功”想成“承诺已登记”把“正在处理”想成“早晚会处理”。如果这篇文章能帮你在设计下一个微服务时别再急着把两个服务用 HTTP 直连在一起而是先问一句“这里是不是应该放一个队列”那它就有价值了。系列后续会继续深入集成模式、消息契约管理、批量消费性能调优和成本优化到时候再一个个聊。