
1. 为什么我会把同步网关链路改成异步消息驱动1.1 原来的架构其实很干净但流量不按套路出牌我接手这个系统的时候架构并不复杂Nginx做负载均衡入口处挂了一个API网关层负责路由、鉴权、限流、灰度再往下是订单、库存、支付这几个核心微服务。服务之间用Feign同步调用链路大概就是网关转给订单服务订单服务调库存服务库存服务再调支付服务一次请求把所有事情干完最后把结果返回给客户端。前期这套结构很舒服。同步模型最好理解开发效率高团队新人上手也快。一个请求进来顺序往下调用日志顺着链路一追到底问题定位非常直观。我还记得系统早期我们一个星期就能上线两三个小功能大家都觉得这套架构再战三年没问题。但等到真实流量上来了问题开始一点点显现。我印象最深的是一次大促前的压测网关线程池最先被打满订单服务的数据库连接池紧接着扛不住整个链路像多米诺骨牌一样塌下去。当时排查了三天根因非常清晰同步调用让大量线程阻塞在下游服务的响应上而下游任何一个服务出现抖动都会顺着调用链传导到网关层。具体来说痛点集中在三个方面。第一是线程资源被无限消耗每个请求在等待下游时都占着一个线程下游慢线程堆积新的请求进不来。第二是依赖链太脆订单服务依赖库存库存依赖支付任何一个服务出问题整条链路成功率都掉而且很难通过简单加机器解决瓶颈可能卡在数据库连接池也可能卡在一个外部接口的响应时间上。第三是大量写操作集中在高峰期用户下单就那么几秒钟但下单之后的库存扣减、优惠券核销、积分发放、短信通知全部同步做完的话用户等待时间被拉长系统的峰值承载能力也浪费在这些用户根本不关心的操作上。1.2 异步化的本质是把“等待”改成“事件通知”那时候我反复在团队里讲一句话用户下单之后其实能接受“已受理”这个状态很多衍生动作完全没必要让用户在线等着。这个想法最终推动了一次架构演进把部分同步调用从请求链路中剥离出去改成基于消息队列的异步消息驱动。异步消息驱动并不是把系统推倒重来它的核心是把链路里的操作分成两类。一类是核心主流程比如订单创建、订单状态更新这些必须同步完成客户端必须在请求周期内拿到明确结果。另一类是衍生动作比如扣库存、发积分、发短信、写物流信息、更新搜索索引这些都可以放进事件里由下游消费者异步处理不阻塞主流程。有一个点特别容易搞混异步化不等于库存扣减不重要了。库存仍然要扣只是“发起扣减”和“扣减完成”解耦了。订单服务只负责发出“订单已创建”的消息库存服务监听这个消息去执行扣减扣减成功后再发“库存扣减完成”的事件往下游走。从调用模型上看原来是一次HTTP请求里串行等三次RPC现在变成一个事件在几毫秒内进入消息总线后续各个消费者按自己的消费逻辑并行处理主流程和衍生动作彻底分开。改造之后网关层也压力骤减。以前网关要把请求转给订单服务然后等订单服务把所有下游都调用完再返回。现在网关的行为变轻了路由、鉴权、限流照旧写请求转给业务服务后业务服务只需要落库订单状态并发送消息然后立刻返回“已受理”。网关线程不再被下游拖死整体吞吐能力上限一下子打开了。1.3 异步化顺带解决了削峰填谷的问题还有一个隐藏收益就是削峰填谷。互联网系统最怕的就是瞬时流量高峰比如秒杀、促销、热点事件带来的流量尖峰如果按照峰值流量部署资源成本极高。消息队列天然就像一个缓冲区瞬时高峰把消息积压下来下游消费者按照自己的能力慢慢消费。我们后续做过一个小实验把同一个促销接口的流量压到原来的三倍同步调用版本直接超时率飙升异步消息驱动版本只是消息积压量涨了一截下游按部就班处理前端用户几乎没有感知。这让我真正体会到为什么很多高并发系统最终都要走消息驱动这条路它不仅仅是性能优化更是系统稳定性的兜底机制。2. 方案设计与中间件选型别急着写代码2.1 异步化改造之前先回答三个问题我见过不少团队一听到“异步化”就把所有接口往MQ里塞结果比原来更乱。异步化确实增加了链路复杂度如果不能在动手前想清楚边界后患无穷。我们在设计阶段强制要求每个候选接口回答三个问题第一用户是否关心这个操作的最终结果比如支付结果查询用户必须立刻知道支付是不是成功了这种绝对不能只发消息就完事必须同步确认。而库存预占、优惠券标记用户在乎的是最终能不能下单成功在请求链路里先做预检查实际扣减放到消费端失败了走补偿这是可以异步化的。第二操作失败后能否通过消息重试补偿还是必须同步返回告警如果失败之后没法自动重试或者重试的成本极高那就要谨慎。有些操作必须严格保证实时性比如账号锁定、安全风控这些不适合异步。第三下游消费者如果挂了业务是否能接受暂缓处理异步化意味着下游故障不会立即反馈到上游但会造成消息积压如果积压过久业务上会产生新的问题比如优惠过期、库存超卖。所以还需要配套制定“积压多久算异常”的监控阈值。我们最终圈定的改造范围是订单创建、库存扣减、积分变动、消息通知、搜索索引更新这一批业务支付结果确认、账务实时冻结、风控拦截这类坚决不动保持同步链路。2.2 消息中间件选型RocketMQ、Kafka、RabbitMQ怎么选市面上的消息中间件可选的不算少中小团队主要就是RocketMQ、Kafka、RabbitMQ三个里面挑也有小部分用Pulsar的但运维成本不低。我们把几个核心维度拉了一个对比表这里也分享出来中间件吞吐能力消息可靠性顺序消息延迟消息死信队列事务消息运维复杂度RocketMQ高高主从同步刷盘支持支持多级延迟支持支持中等配套完善Kafka极高高ISR机制分区内有序需自研需自研有但弱中等依赖ZooKeeper或KRaftRabbitMQ中等高单队列有序支持支持需插件低最终我们选了RocketMQ。原因很直接Kafka的吞吐量确实最强但它的定位偏日志管道很多面向业务的消息语义需要自己搭比如延迟消息、事务消息、死信队列RocketMQ基本上开箱即用。RabbitMQ更轻运维也简单但消息量一旦涨到分区级别队列性能下滑和集群稳定性都不如RocketMQ。考虑到我们要做订单、库存这类对事务和顺序有明确要求的业务场景RocketMQ的综合匹配度最高。2.3 Topic设计是架构的地基别把一锅粥放在一个队列里Topic设计是我在异步化改造里吃过亏最多的地方也是被大多数团队忽略的地方。Topic划分太粗各种业务事件混在一起消费者要做大量无谓的消息过滤消费端代码写起来非常臃肿。划分太细Topic数量爆炸运维和管理成本陡增。我们最终确定了“按业务域划分按事件语义命名”的原则并且明确区分命令消息和事件消息两种类型命令消息比如“扣减库存”“核销优惠券”目标是明确通知某个服务执行某个动作。事件消息比如“订单已创建”“支付已完成”表示某个事实已经发生生产方不关心谁来消费。这两种消息的语义差别很大发送方和消费方的协作契约也不同。命令消息要求消费方必须执行成功失败要重试或告警事件消息更像广播下游自己决定如何处理。实际Topic命名方面我们用下单链路举个例子。订单服务发一个“ORDER_CREATED”事件库存服务监听后扣库存扣完再发一个“INVENTORY_DEDUCTED”事件支付服务监听后发起支付支付完成发“PAYMENT_COMPLETED”。每个Topic的事件Payload里带的是业务事实本身比如orderId、userId、商品列表、金额而不是“请调用某某接口”这种命令式参数。这样设计的好处是下游消费者可以灵活决定自己的处理逻辑多个消费组监听同一个Topic各做各的事情互不干扰。3. 实操记录从订单服务开始做第一个异步化改造3.1 生产端改造先保证主流程落库再发消息我拿最典型的下单链路做讲解。改造之前的Controller里直接调库存、优惠券、积分三个服务改造之后变成先保存订单状态为“已创建”然后向Topic发送“订单创建成功”事件随即返回客户端。这里有一个极其关键的细节生产端到底应该先落库再发消息还是先发消息再落库很多人初学时会直接先发消息但这样做风险很高如果消息发成功了业务事务回滚了下游消费者会拿着一个不存在的订单去扣库存出现脏数据。反过来如果先落库再发消息消息发送失败了也能靠补偿机制补发。我们采用的方式是“本地消息表 定时任务补偿”。在订单表旁边建一张消息记录表业务事务提交成功后在同一条数据库事务里写一条初始化的消息记录消息发送成功后再更新这条记录为已发送状态。如果发送失败定时任务会扫描超时的未发送记录重新推送从而保证消息不丢。生产端核心逻辑如下Transactional public OrderVO createOrder(CreateOrderRequest request) { // 1. 核心订单数据落库状态 CREATED Order order orderRepository.save(buildOrder(request)); // 2. 同事务记录一条业务消息用于可靠发送 LocalMessage message LocalMessage.of(order.getId(), ORDER_CREATED, buildEventPayload(order)); localMessageRepository.save(message); // 3. 事务提交后异步发送并刷新消息状态 TransactionSynchronizationManager.registerSynchronization( new TransactionSynchronization() { Override public void afterCommit() { SendResult result messageProducer.send( order_event, ORDER_CREATED_ order.getId(), order.getId(), message.getPayload()); if (result.getStatus() SendStatus.SEND_OK) { localMessageRepository.markSent(message.getId()); } } }); return OrderVO.from(order); }这里有几个细节值得注意。发送消息时我把orderId作为消息的key传进去了这是为了后续做顺序消息同一个订单的消息可以保证进入同一个分区消费端处理时会方便很多。另外发送逻辑放在事务afterCommit里而不是直接在业务代码中同步发送原因是避免数据库事务还没提交时消息已经发出去了消费端查询不到订单数据导致空处理。3.2 本地消息表和事务消息我们为什么选前者很多团队会直接用RocketMQ的事务消息机制。这个方案的原理是先发一条半消息业务本地事务提交成功后把半消息置为可投递状态业务失败就删除半消息。如果中间长时间没有收到确认MQ会主动回查本地事务的结果决定是否投递。我们当时评估过事务消息但最终选了本地消息表方案。原因有两个一是我们团队对数据库非常熟悉本地消息表可以随时查看消息发送状态可观测性更强出了问题直接查数据库就行二是事务消息对TPS有一定损耗回查机制的实现也容易出错一旦回查状态逻辑写反后果很严重。本地消息表虽然多了一张表但对于核心链路的稳定性保障完全值得。3.3 消费端改造幂等是异步架构的命根子消费端相比生产端要谨慎得多。消息队列的投递语义是“至少一次”也就是说在极端情况下消息可能会被重复投递。消费端如果没做幂等就会出现重复扣库存、重复发短信这种事故。我第一个消费端就踩了这个坑后来总结出一套标准模板public void onOrderCreated(MessageExt message) { String orderId message.getKeys(); String dedupKey dedup:order_created: orderId; // 1. 用 Redis SETNX 做消费幂等 Boolean first redisTemplate.opsForValue() .setIfAbsent(dedupKey, 1, 24, TimeUnit.HOURS); if (!Boolean.TRUE.equals(first)) { // 重复消息直接返回 return; } try { // 2. 核心业务逻辑 inventoryService.deductStock(orderId); } catch (Exception e) { // 3. 消费失败抛出异常让MQ根据重试策略重新投递 log.error(deduct stock failed, orderId{}, orderId, e); throw e; } }这里有一个非常多人写错的细节业务逻辑执行失败后要不要立刻删除Redis里的幂等标记我的答案是绝不删除。消息消费失败后MQ会重试重试时第一步SETNX如果能设置成功消费逻辑就能重新执行一遍这样才有补偿的机会。如果你把幂等标记删掉了重试时会被误判为重复消息直接跳过业务一辈子都不会成功。当然Redis幂等方案也有局限如果业务处理耗时较长超过幂等key的过期时间就会出现同一消息重复执行。我们最终在生产环境做了一个升级对于订单、库存这类核心业务用数据库唯一索引做持久化幂等商业上更加可靠。CREATE TABLE msg_consume_log ( id BIGINT PRIMARY KEY AUTO_INCREMENT, msg_key VARCHAR(64) NOT NULL, topic VARCHAR(64) NOT NULL, consumer_group VARCHAR(64) NOT NULL, consume_status TINYINT NOT NULL DEFAULT 0, consume_time DATETIME DEFAULT NULL, UNIQUE KEY uk_msg_consume (msg_key, topic, consumer_group) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;消费端开启事务先insert这条消费记录唯一键冲突说明是重复消息直接忽略。如果插入成功才执行业务逻辑业务逻辑和消费记录插入在同一个本地事务中要么一起成功要么一起回滚。这套方案上线后重复消费导致的业务问题基本绝迹。4. 上线后踩过的坑以及排查实录4.1 重复消费搞出来的“短信轰炸”改造后第一周就碰到一次事故。用户反馈收到了重复的优惠券领取短信而且每个人收到的次数还不一样有人收到两条有人收到五条。排查之后发现是两个问题叠加出来的。第一起初我把幂等key的过期时间设置成了30分钟但有一个下游消费者因为外部短信接口抖动一条消息积压了超过30分钟才开始处理中间MQ又重投了几次每次幂等key都重建了业务逻辑就执行了多次。第二消费成功逻辑里还写了一段“处理完成后主动删除幂等key”的代码这两个问题叠一起重复消费被彻底放大了。修复方案就是前面讲的核心业务的幂等从Redis换成数据库唯一索引过期时间不再依赖Redis的TTL而是永久保留记录。同时把消费成功日志里自动删除幂等key的代码全部下线。4.2 顺序消息库存扣减千万别跑到订单创建前面入手异步化之后业务又开始涉足状态流转类场景。订单从“已创建”到“已支付”再到“已发货”状态变化是有严格顺序的如果“已支付”的消息先被消费而“已创建”的消息还在队列里排队就会发生状态倒置业务上完全不可接受。RocketMQ对顺序消息的支持是“局部顺序”也就是同一个业务key的所有消息会被路由到同一个队列消费者对该队列单线程处理从而保证顺序。关键实现有两个点。生产端要使用MessageQueueSelector根据业务key做哈希取模选队列SendResult sendResult producer.send(message, new MessageQueueSelector() { Override public MessageQueue select(ListMessageQueue mqs, Message msg, Object arg) { String orderId (String) arg; int index Math.abs(orderId.hashCode()) % mqs.size(); return mqs.get(index); } }, orderId);消费端要保持每个队列单线程消费不要在消费方法里再引入线程池异步处理否则顺序会乱。很多初级的坑是生产者配好了按订单哈希路由但消费者内部为了追求速度快把消息丢进线程池并发处理结果顺序完全失去了意义。还有一点要注意顺序消息Topic的分区数尽量不要动态扩容。如果队列数从4扩容到8原来hash到队列0的订单一方面另一批新的消息可能被hash到别的队列同一订单的后续消息跑到另一个队列去了疫情期间的乱序问题会在扩容窗口集中爆发。所以线上对顺序性要求高的Topic扩容一定安排在业务低谷而且要有消息对账机制兜底。4.3 消费积压把数据库连接池打爆了一次大促期间订单量暴增某个下游消费者消费能力跟不上消息积压了几百万条。当时没有第一时间发现直到数据库连接池报警被打爆才顺着监控反查到消息积压。排查过程是这样打开RocketMQ控制台看到某个Topic的积压数量持续上涨消费组在线状态但消费TPS一直上不去。再看消费日志大量消费线程全部阻塞在调用外部优惠券接口的位置这个外部接口单次响应已经超过5秒。消费者线程处理一条消息需要5秒消息排队自然越来越长。这种问题的标准解法是把消费端拆成两层第一层快速接收消息只做数据入库和去重然后立刻确认消费。第二层通过线程池异步去调用那些慢速外部服务不再阻塞消费线程。但要注意拆分之后业务成功确认就变复杂了第一层确认了不代表第二层成功最终还是要靠“业务状态标记 定时对账”来兜底。我后来养成的习惯是消费端必须监控两个基础指标单条消息处理耗时和消费TPS。平均处理耗时一旦超过500毫秒必须人工介入调查潜在的问题越早暴露越好。4.4 死信和重试别让一条坏消息堵住整个队列消息消费失败后如果不断重试会一直占据消费线程把后面的消息全部挡住。RocketMQ的默认策略是消费失败后重试16次间隔时间从10秒逐步扩大到2小时16次之后进入死信队列。这个策略最大的问题就是间隔时间太长了。如果有一条消息因为代码bug会稳定失败它会在消费端反复重试十几个小时期间把这个消息所在分区后面的所有消息全部堵住整个队列的消费进度停滞不前。我们后面给不同的业务配置了不同的重试策略。订单类核心链路用快速重试最多重试3次间隔分别是5秒、30秒、5分钟三次都失败直接进死信队列。这样即使有问题也能快速暴露不会长时间阻塞队列。死信队列也要有人盯。我们写了一个定时任务定期扫描死信队列里最近1小时的新消息自动推送到告警群同时把消息内容刷进一张死信记录表方便人工查看原因和手动重放。4.5 消息轨迹与全链路追踪异步排查的救命稻草异步化改造最大的痛点之一就是排障模式的改变。以前同步调用一条日志链路直接追踪到底现在消息被解耦之后你从网关看到一个订单号但后续消息在哪个节点消费、处理耗时多少、最终成功还是失败如果不在同一个日志链路里排查起来非常痛苦。我们的解法是给每个请求和每条消息绑定统一的traceId。网关入口生成traceId放进HTTP Header业务服务在处理时把traceId取出来透传发消息的时候把traceId写进消息属性消费者消费时从消息属性取出traceId打印到自己的业务日志里。这样只要拿到一个订单号在整个日志平台上搜索traceId就能把网关、订单服务、库存服务的日志按时间线全部串起来。异步链路的排查效率提升非常明显。我强烈建议把全链路追踪放在改造初期就做后面再补会很麻烦。4.6 监控指标与告警规则异步系统的眼睛异步化改造完成之后团队对“系统健康”的判断方式也要跟着变。以前只要盯着接口错误率和响应时间后来我发现还需要增加几个消息维度的指标。生产和消费QPS的差值如果差值持续扩大说明消费能力跟不上生产速度。每个消费组的积压数量和积压消息的最大年龄积压年龄比积压数量更重要积压10万条但如果1分钟能消费完其实不可怕积压1000条但是年龄已经1小时反而说明消费已经卡死了。还有消费失败率和重试次数分布失败率异常升高往往意味着下游依赖服务出了问题。我们把这几项指标都加到了告警规则里阈值是消息积压量超过5万条或者积压消息最大年龄超过10分钟。上线这套监控后再没出现过大促期间消息积压到数据库连接池被打爆的情况。5. 容量规划与团队协作异步化落地后我的一点复盘5.1 消息集群不是无限的缓冲池容量还是要算的关于异步消息驱动很多人有一个误解觉得消息队列是万能的“大水池”可以无限吞流量。但MQ集群本身也有网络带宽、磁盘IO、存储空间的上限它只是比业务服务的缓冲能力强很多不是无限强。容量规划上我们按大促峰值下单量的三倍作为生产端峰值来计算。单个broker节点实测能稳定扛住每秒8000条消息的生产和消费预估大促峰值是每秒1.5万条那就至少准备3个broker节点留出冗余。磁盘方面按消息保留3天来计算单条业务消息平均大小约1KB日均消息量在亿级别估算一天大概100GB3天就是300GB再加上副本冗余每台broker准备500GB磁盘上线后基本稳定。5.2 团队协作方式也跟着变了异步化改造不只是一个纯技术问题还会改变团队内部的协作模式。以前订单服务出问题查日志链路直接锁死是哪个下游服务现在消息解耦订单服务只关心自己发没发消息库存服务要关心自己消费成不成功两边必须对消息格式有清晰的约定。我们后来整理了一份内部的消息规范文档内容包括事件命名的语法、每个Topic的用途和负责人、消息Payload的JSON字段定义、消费者组的命名规范、幂等key的定义规则。这份文档成了团队后续迭代和新人培训的重要素材。没有这份规范的话团队越大越容易乱。5.3 异步化改造的真实收益和一句忠告这次改造上线后系统峰值吞吐能力提升明显核心下单接口的P99响应时间从800毫秒降到120毫秒网关层线程池从频繁报警变得基本稳定。最让我满意的不是性能数据而是系统面对故障时表现“软”了一个下游服务抖动不再像以前那样立刻传导到整条链路而是通过消息积压、消费降速慢慢体现出来给团队留下了处理时间。如果让我重新做一次我会一开始就安排两件事第一消息规范文档在动手写代码之前先定稿第二把消息轨迹和全链路追踪优先排期而不是等出了问题再补。异步化本身不是银弹它只是把同步问题的复杂度转移到了消息中间件和分布式一致性上用好了是解放用不好就是另一个运维噩梦。最后再给一个建议如果你的系统还在同步调用阶段不要一上来就全量异步化。挑一个链路最长、对实时性要求最低的业务场景做试点把监控、幂等、补偿这一整套闭环跑通再逐步推广。这条路我走了一遍方向是对的剩下就看你的系统要不要做这个选择了。