资讯详情

资讯详情

从同步网关到异步消息驱动:高并发系统架构演化实践

1. 当初为什么要动这套架构网关同步模式的真实痛点1.1 旧系统长什么样接入网关加同步转发的经典组合我在上一家公司负责的是一个电商中台核心系统的架构维护业务形态比较典型用户端是微信小程序和App后端有会员、订单、库存、营销等多个微服务。在很长一段时间里我们的对外接入层用的是API网关请求统一打到网关再做鉴权、限流、参数校验然后同步转发到下游的各业务服务。这套架构放到2018年那个节点应对日活几百万的读写请求是没问题的。网关用的是OpenResty那套基于Nginx和Lua脚本做的二次开发限流、鉴权、灰度路由这些能力都写在Lua模块里下游服务的注册发现走的是Consul服务间的调用是Feign走HTTP同步接口。从调用链路的视角看一次完整的用户下单是这样的客户端请求到达网关网关根据URL前缀路由到订单服务的某个接口订单服务在接口内部同步调用库存服务扣减库存调用会员服务锁定优惠券调用支付服务生成预支付单所有环节都是同步等待返回全部成功才算下单成功。这里有一个非常典型的特点同步调用链路的成功与否取决于链条上最慢的那个环节。我们当时的接口平均响应时间在300毫秒左右看起来数字不算难看但问题是这种系统没有任何抗突发流量的能力。每年的双11大促前团队的主要工作就是在网关层加限流阈值、扩容下游服务实例说白了就是用资源硬扛。网关层的限流配置写得再精细也只是把超出处理能力的请求拒之门外系统的弹性完全靠机器数量堆出来的。1.2 高并发冲击下暴露出来的真实问题真正逼我们动手改造的是三次线上事故的叠加。第一次是营销活动秒杀场景。运营搞了一场限量5000件商品的秒杀活动刚开始的瞬间网关层的QPS瞬间从平时的2000出头打到了4万多。限流规则起了一定的保护作用但下游的订单服务照样被同步调用打挂因为订单服务对库存服务的扣减接口是同步等待的库存服务扛不住的时候订单服务的Tomcat线程池被占满后续请求全部阻塞在排队队列里最终引起雪崩。那次事故的直接损失是活动页面白屏了40分钟技术团队被业务部门点名批评复盘结论写的是网关限流策略不够精细化但实际上我们都清楚根子出在同步调用的耦合方式上。第二次是商家后台的批量导入功能。运营人员通过后台批量导入10万条商品数据导入服务需要逐条调用审核服务的同步接口一条数据审核平均需要50毫秒10万条串行跑完要将近一个小时。中间只要有一条数据触发了审核服务的异常整个导入任务中断前面审核通过的几万条数据状态全部悬空需要人工对账修复。这类问题用同步调用的方式几乎无解不是换个网关、加个熔断就能解决的。第三次是数据洪峰把数据库连接池打穿。某个合作渠道的回调通知在深夜集中推送过来通知服务同步调用订单服务更新订单状态订单服务再同步查询数据库做状态校验。那一波回调的量其实不算大但打了我们一个措手不及数据库连接池瞬间被占满主库的CPU直接飙到100%导致整个核心交易链路全部不可用。事后分析发现问题不在于某一台机器的性能而在于每个需要更新订单状态的环节都同步占用着数据库连接一个慢查询影响整条链路。这三起事故让我意识到一件事同步调用模型本质上是一种时间的强耦合。上游和下游共享同一个线程生命周期共享同一个响应超时时间共享同一个可用性命运。当系统规模跨过某个临界点之后继续在网关和同步框架上修修补补投入产出比会越来越低。1.3 什么样的信号出现就该考虑架构演化了这里想给还在用同步网关架构的团队一些判断依据。不是所有系统都需要演化为异步消息驱动架构小规模业务用同步调用反而更简单、更可控。但当下面几个信号集中出现的时候就值得认真评估演化的必要性了。第一个信号是高峰流量下的系统可用性依赖扩容。如果你的团队在每次大促前都要靠预估流量、提前扩容来保系统稳定扩容之后线上依然频繁出现接口超时说明系统本身的弹性不足而不是资源不够。第二个信号是业务链条中存在明显的长尾慢操作。比如文件处理、批量导入、外部API对接、短信和邮件通知这些操作如果它们被塞进同步接口链路里会导致接口响应时间被严重拖长。用户点击一个按钮前端转圈转半天体验极差。第三个信号是上下游团队之间的交付协调成本越来越高。订单服务改了字段库存服务要跟着改接口联调整整改一两天营销团队希望在用户下单时做实时推荐但又不想影响主链路响应时间。这种跨团队的接口耦合本质上就是同步调用模型带来的组织协同成本。第四个信号是同样的逻辑被同一个上游以不同方式反复调用而且调用方之间对成功标准的定义不一致。典型例子是订单状态变更通知有的调用方希望同步拿到结果有的希望异步接收通知有的调用失败后会自动重试有的不会。这部分逻辑如果全部由上游各自实现迟早会因为某次异常重试导致数据重复或丢失。我们当时对照这些信号逐一排查发现四个信号全中。于是团队开始正式讨论异步消息驱动架构的改造方案。2. 方案选型为什么目标定在异步消息驱动2.1 先排除掉的那些路线真正讨论方案的时候团队内部提过好几条路线这里简单说说为什么最终没有选。第一条路线是继续优化网关和同步链路。具体做法包括网关层精细化的限流熔断策略、下沉公共逻辑到网关、加多级缓存降低数据库压力、服务间调用全部引入超时控制和隔离线程池。这条路线我们其实已经做了不少在改造前网关的限流策略已经细化到按用户维度、按接口维度分别配置缓存命中率也优化到了85%以上。但它解决不了本质问题同步调用链路的可用性仍然取决于最薄弱的下游服务缓存只能挡住读请求写请求的链路依然是同步阻塞的。第二条路线是把网关替换成Service Mesh方案用Sidecar做服务间通信治理。我们当时调研过Istio但在2019年那个阶段Service Mesh在性能损耗和运维复杂度上的代价还比较高团队也没有足够的容器化基础设施支撑它落地。从投入产出比来看这条路线适合基础设施非常完善的大厂我们的现状是连Kubernetes集群都还在迁移过程中贸然上Service Mesh风险太大。第三条路线是引入响应式编程框架用Reactor模型改造服务间的调用方式。WebFlux、RxJava这套方案在提升单机并发能力上确实有效但它对业务代码的侵入性特别强团队已有的同步代码改造成响应式风格工作量不亚于重写一套系统。加上响应式编程的排查链路经验不足出了问题很难快速定位对中小规模团队来说学习成本和运维成本都偏高。综合对比后团队最终达成了共识核心改造思路不是把同步调用优化成异步调用那么简单而是把同步依赖改造成事件驱动。具体落地手段就是引入消息中间件把服务间的强依赖关系转化为基于消息的弱依赖关系。网关层继续保留但网关的职责收敛为请求接入、鉴权、限流和路由分发业务处理逻辑全部异步化通过消息驱动下游服务完成各自的职责。2.2 消息中间件选型我们为什么选了RocketMQ方案确定后团队内部关于选型也讨论过好几轮。市面上主流的消息中间件无非是Kafka、RocketMQ、RabbitMQ这三个我们最终选了RocketMQ这个决策的过程值得分享。先说说为什么排除RabbitMQ。RabbitMQ的优点是轻量、部署简单、社区成熟基于Erlang语言开发的性能表现也不错。但它在以下几个方面不能满足我们的需求一是消息堆积能力不够强RabbitMQ的队列深度上去以后内存占用飙升我们的业务场景里秒杀活动的峰值消息量非常大可能需要短时间堆积上千万条消息二是它的顺序消息支持不够方便虽然可以通过single active consumer实现但整体使用成本偏高三是运维生态相对独立公司内部对它的监控告警和管理经验积累不足。然后是Kafka和RocketMQ的对比。Kafka的吞吐量确实比RocketMQ高一个量级在大数据日志采集领域是绝对的主流。但我们当时的核心诉求并不是极致的吞吐量而是业务消息的可靠性、事务消息和延迟消息这些功能。RocketMQ在功能层面上完胜事务消息可以很好地支撑本地事务与消息发送的原子性问题延迟消息直接支持多个延迟级别不需要自己实现定时调度逻辑消息消费支持tag过滤可以按业务事件类型精确订阅官方控制台提供了完整的管理和监控界面落地门槛比Kafka低很多。最关键的因素是RocketMQ的社区在阿里巴巴内部经过了双11的大规模验证很多我们即将踩的坑它都已经踩过了。对于一个需要快速落地、团队人手有限的场景来说选择一个人有验证过的成熟方案比追求技术指标的极致更重要。最终我们选定RocketMQ 4.5版本部署模式是四主四从的集群每个主节点配一个从节点主从之间采用异步复制模式。之所以不全量同步复制是为了避免复制链路拖慢主节点的写入性能在业务容忍极少量消息丢失的前提下异步复制的性价比最高。2.3 目标架构网关瘦身加事件驱动改造后的目标架构我画了一个很简单的逻辑模型用户请求经过网关完成最基础的接入管控后网关直接响应客户端已受理同时把业务事件发布到消息队列之后的处理过程全部由消费者异步完成。拿用户下单这个核心链路举例改造前的处理逻辑是同步的网关转发到订单服务订单服务调库存扣减库存、调会员服务核销优惠券、调支付服务生成预支付单全部同步完成才返回。改造后的逻辑变成网关收到下单请求做基础校验后生成一个订单创建事件发布到RocketMQ的trade_order_create主题然后立刻返回客户端下单请求已受理。订单服务作为消费者异步处理这个事件先做业务校验发送预支付单生成请求到支付渠道支付成功后发出order_paid事件库存服务监听order_paid事件扣减库存会员服务监听order_paid事件核销优惠券。整个过程对客户端来说响应时间是网关内的处理时间一般在20毫秒以内而业务全流程的处理时间从原来的300毫秒变成了可能1秒以上但用户没有感知到变慢因为前端拿到已受理就开始轮询订单状态。这个模式其实是很多互联网系统在用的效果一致性方案对用户来说感知不到异步处理的延迟只要轮询接口能及时反馈状态变化即可。网关在这个架构里的角色发生了质的变化它不再是一个流量转发和业务聚合的中间层而是一个纯粹的接入控制面和事件发布面。限流规则不再需要精细到每个下游服务的负载情况只需要保护网关自身和消息队列写入的稳定性即可。下游服务的任何故障都不会直接拖垮网关因为网关不等待下游的处理结果消息队列天然起到了削峰填谷和故障隔离的双重作用。3. 灰度迁移实操从双写到切流3.1 第一步把写操作改成异步化架构目标定了落地过程不能搞一刀切。我们用了将近三个月的时间分阶段完成了从同步网关到异步消息驱动的过渡。第一个阶段是双写验证。在核心链路上同步调用照常走但订单服务在处理完同步逻辑后额外把订单创建事件发送到RocketMQ用于验证消息生产、消费链路的畅通和数据的完整性。双写阶段持续了两周期间我们做了大量的消息内容校验工作写了一个专门的对比程序定时拉取消息队列中的事件和数据库中的订单记录做比对确保消息内容没有丢字段、没有序列化异常。第二个阶段是影子消费验证。消息队列里的订单创建事件已经能被消费者正常处理了但消费者这边只做数据流转不真正影响核心业务结果。比如库存服务消费order_paid事件后只是记录一条日志应该扣减库存xx件不实际扣减。这个阶段核心验证的是消息处理的幂等性和延迟指标确保消费者处理消息的耗时不会出现明显的尖刺。第三个阶段才是真正切换。我们按业务线逐个切流量先用占比5%的用户切到异步链路观察几天看有没有异常再逐步放大到10%、30%、50%、100%。整个过程用了一个多月每次放量都配合监控指标的观测至少有订单创建率、支付成功率、消息堆积量、消费者处理耗时这几个核心指标在盯着任何一个指标出现异常都会立刻回滚。这里有一个核心经验分享切换顺序一定要先从低风险业务开始。我们第一个切换的是订单创建这个场景因为订单创建不像支付和库存扣减那样涉及资金和库存的强一致性要求即使异步处理出问题最多是订单状态长时间不更新不会造成资金损失。库存扣减和支付回调我们放到了最后一批切换因为这两个环节对一致性的要求最严格。3.2 消息生产与消费端的关键参数这块是实操性最强的部分我们把RocketMQ生产端和消费端的核心参数都过一遍。生产端我们用的是RocketMQ的DefaultMQProducer核心参数有这几个sendMsgTimeout默认3000毫秒我们生产环境调到了5000毫秒。原因是生产端发消息偶尔会因为Broker刷盘耗时或者主从复制存在小的波动超时时间太短容易误判发送失败。retryTimesWhenSendFailed同步发送失败的重试次数默认是2我们设置成3。这里要注意的是RocketMQ的重试机制在同步发送场景下可能会产生重复消息所以消费端必须做幂等处理。compressMsgBodyOverHowmuch消息体超过4KB自动压缩我们保持默认不额外调整。生产端的发送方式我们全部采用同步发送没有用异步发送或者单向发送。原因很简单同步发送能找到明确的发送结果出错时可以在业务代码里做补偿处理异步发送虽然有吞吐量优势但错误处理逻辑要写在回调里代码可维护性差一些。消息发送失败时我们的兜底策略是记录一条失败日志到本地数据库由定时任务扫描补偿把失败的消息重新投递到RocketMQ。这个补偿机制保证了消息不丢失是异步化改造里的最后一道保险。消费端用的是DefaultMQPushConsumer核心参数更值得关注consumeThreadMin和consumeThreadMax消费者线程数我们设置的是20到60之间。这个线程数的设置不是越大越好因为消费速度取决于下游业务的处理能力线程数太大反而会因为频繁上下文切换降低吞吐量。consumeMessageBatchMaxSize每次消费的消息数量默认是1意味着每次只拉取一条消息处理。我们根据业务处理耗时设置到了8。这个值需要配合消费耗时来调整单个消息处理耗时在50毫秒以内的场景批量消费能显著提升吞吐量。consumeTimeout单条消息消费超时时间默认15分钟我们设置成30分钟。如果消费者的业务处理逻辑涉及外部接口调用有时候单条消息处理时间可能超过默认超时导致RocketMQ判定消费失败触发重投递。消费端最关键的配置是消息重试策略。RocketMQ默认的重试次数是16次重试间隔按指数递增。我们的经验是不要使用默认的重试次数而是要结合业务情况收敛重试次数。对于数据类消息16次重试堆在队列里会造成消息大量积压影响后续消息的处理时效。我们把重试次数调整到5次超过重试上限的消息进入死信队列由人工介入处理。这样既保证了消息不会轻易丢失又避免了无限重试导致的队列堵塞。3.3 切流过程中的坑幂等、乱序与延迟迁移过程中我们踩了一堆坑有些坑不踩一遍根本想不到这里选几个典型的分享。第一个坑是消息重复消费导致数据被重复处理。RocketMQ的消息投递语义是at least once也就是说消息可能被重复投递。第一次切流的时候我们订单服务的消息消费者在处理创建订单事件时没有做严格的幂等控制结果上游某次网络闪断触发了消息重试同一个订单事件被消费了两次数据库里出现了两条一模一样的订单记录。排查了半天才定位到这个原因后来在消费者入口统一加了幂等处理方式很简单在本地业务表里为每条消息建一个唯一索引消息的msgId作为唯一键处理前先插入消息记录插入成功才执行核心业务逻辑插入冲突说明消息已处理过直接跳过。第二个坑是部分场景下消息乱序导致的数据错乱。典型场景是订单状态变更一笔订单先收到已支付事件又收到已取消事件正常情况下应该先取消后支付但因为消息是分布在不同队列并行消费的顺序无法保证最终数据库里订单状态被错误地更新成了已支付。处理方案是对强顺序要求的消息类型单独使用RocketMQ的顺序队列。RocketMQ的普通消息队列只能保证单队列内的顺序我们在应用层做了一个订单号哈希分桶的策略把同一个订单的事件都路由到同一个消息队列再配合单队列单线程消费的方式保证订单侧的消息严格有序。第三个坑是消息延迟造成用户体验下降。异步化之后用户下单后客户端可以立即收到已受理的响应但如果消息的消费延迟超过用户的心理预期用户刷新订单列表看不到新订单就会产生投诉。我们当时就遇到过这种情况某个凌晨RocketMQ的一个Broker节点磁盘告警消息写入变慢消费端拉取消息的延迟从正常的上百毫秒飙升到了十几秒用户在App里下单后迟迟看不到订单记录。后来在监控上加了消费延迟时间的告警线超过5秒就开始告警同时针对关键链路的订单查询接口做了兜底逻辑查询不到异步创建的订单时直接查RocketMQ中待处理的消息尽可能让用户感知不到异步处理的延迟。4. 演化完成后的效果和新挑战4.1 压测与线上数据对比改造全部落地之后我们对系统做了一次完整的压测也把线上数据拉出来做了对比。这里把数据放出来供参考。网关层的核心指标改善最直观。改造前网关到订单服务的同步转发链路P99响应时间在280毫秒左右改造后网关发布消息到RocketMQ的P99响应时间稳定在15毫秒以内。网关层几乎不再受下游服务故障的影响下游服务整体宕机5分钟网关层都不会出现请求排队和超时。系统整体吞吐能力的变化更大。改造前订单创建链路的压测极限是5000 QPS瓶颈在订单服务同步调用库存服务时的数据库连接池上继续加压就会雪崩。改造后同样规格的订单服务可以支撑的峰值写入量提升到了30000 QPS以上而且系统表现是平滑的QPS再高也只是消息积压量增加下游服务处理不过来时消费速度自动降速不会出现服务不可用的情况。资源成本上改造前大促时需要给订单服务扩容到40个实例改造后同样流量下只需要24个实例。成本节省接近40%而且大促期间的扩缩容压力小了很多不再需要根据预估流量提前一周大规模扩容只需要保证RocketMQ集群的Broker节点数量充足下游消费者实例是水平可扩展的流量上来后动态伸缩即可。4.2 异步化带来的新麻烦排查链路与分布式事务异步架构也不是银弹它解决了一部分问题必然带来另一部分新问题。这里分享我们实际遇到的三个新麻烦。第一个新麻烦是排查链路变长了。同步调用时代一次请求的全链路信息可以通过traceId串联起来从网关入口到下游服务每条日志里都有traceId查一次日志就能看到整个调用链。异步化之后一条业务请求被拆成了多个事件每个事件的消费端都会生成新的上下文单纯靠traceId已经没法串联全链路了。我们的解决方案是引入OpenTracing规范把生产消息时的traceId和spanId透传到消息体里消费者在消费消息时恢复上下文重新建立新的span把父子关系通过消息体传递。这样排查问题时就能把整条异步链路的调用关系串起来。第二个新麻烦是分布式事务的复杂度变高了。同步调用的时候我们可以用Seata这类分布式事务框架做全局事务控制保证跨服务调用的强一致性。但异步化之后整个链路变成了最终一致性模型订单服务接收创建事件成功后消息已经落库如果再发现订单信息有误没法简单地回滚已经发出的消息因为下游的库存服务可能已经在消费这个事件了。我们比较务实的处理方式是区分业务可靠性等级采用不同的设计模式。对于资金相关的场景比如支付成功后的结算处理核心逻辑是消费端主动查询、对账补偿而不是依赖消息回滚。每个消费端在处理消息时都会先查询上游的权威数据源做二次确认确保数据没有问题才执行写操作同时在每天凌晨跑一次全量对账任务把消息处理的最终一致性兜住。第三个新麻烦是消息堆积的治理从罕见的事故变成了日常的运维。异步架构下消息堆积是常态重点在于区分哪些堆积是可以接受的哪些需要介入处理。我们建立了一套消息堆积的分级告警机制按主题维度统计消费延迟时间延迟超过1分钟告警超过5分钟升级到值班负责人超过30分钟触发紧急响应。同时针对不同主题设置了不同的堆积容忍度日志类消息可以堆积一小时也无所谓但订单状态类消息不能超过5秒。4.3 给准备做异步化改造的团队一些经验清单整个改造做下来如果要浓缩成几条经验我会把这些列在最前面第一异步化改造前先要把幂等这件事想清楚并提前在业务代码里落地好幂等控制。消息队列本身就带有at least once的语义配合生产端的发送重试一定会产生重复消息。幂等不是消费者自己的事而是整个消息链路都需要考虑的消息的msgId生成规则、消息体里业务唯一键的设计、消费端幂等表的结构这些细节在建链路的初期就要定好。第二不要把异步化改造和业务功能迭代混在一起做。我们当时吃过这个亏有一段时间一边改消息队列一边接新的营销活动需求结果两边互相干扰出了问题都不知道是改造引入的还是新功能引入的。后来定了一个规矩改造期间业务功能冻结所有需求排队等改造完成再上线。这是很痛苦的决定但保证了改造期间的线上稳定性。第三消息中间件不是越先进越好而是要选择团队能驾驭的方案。我们现在有同事会问我为什么不直接用Kafka我的回答是如果我们的核心诉求是日志采集和流计算那Kafka无疑是最合适的但我们的核心诉求是业务事件驱动和分布式事务RocketMQ在功能匹配度上更合适。选型这件事技术先进性排在后面团队掌握程度和运维成本排在前面。第四一定要提前建设消息链路的可观测性。至少包括三个方面消息生产端的发送耗时、发送成功率、失败原因分布消息Broker的Topic积压量、消费延迟、Broker节点磁盘和CPU水位消息消费端的消费速率、消费失败次数、重试队列长度。这些指标在改造前期就要接入监控系统否则等到切流的时候再补就来不及了。第五用户侧的体验兜底设计。异步化之后用户感知到的部分操作变慢了需要在产品层面设计好对应的反馈机制。我们做的是在订单列表和订单详情页增加了一个处理中的状态同时前端会轮询订单状态接口拿到最终结果后自动刷新页面。这部分如果不在改造时一并规划用户就会觉得系统变卡了本质上不是系统卡而是用户等待的反馈链路变了。根据我个人在实际操作中的体会异步消息驱动架构最适合的改造时机不是系统已经稳如泰山的时候而是团队已经能明显感知到同步调用的痛点、又还没有被历史包袱拖垮的阶段。一旦跨过那个临界点同步调用带来的问题会越来越频繁地暴露在线上事故里到那时改造的成本会比现在高好几倍。如果你所在的团队正在犹豫要不要走这条路我的建议是先挑一条非核心业务链路做一次小范围的异步化验证用数据说话比在会议室里争论一整天都管用。
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →