ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

IM消息转发子服务:从拆分到高并发落地的完整复盘

IM消息转发子服务:从拆分到高并发落地的完整复盘 做IM后端这几年我最大的一个体会是消息转发这层看着只是把A的话传给B一旦上了规模、上了多端、上了群聊它就成了整个系统里最容易翻车的部位。早期我们第一版IM甚至没有独立的“消息转发子服务”代码就塞在业务服务里结果一回压测就把问题全暴露了。这篇就围绕这个子服务把我从拆分、设计、落地到线上排查的完整过程做一个复盘希望能给正在搭IM或准备搞IM的团队一点参考。1. 拆服务之前先说清楚消息转发到底承担了什么职责1.1 一开始我把转发逻辑写在业务服务里吃了三个亏最早版本的消息链路很直接客户端通过长连接网关把消息体交到消息业务服务业务服务做权限校验、会话校验、写MySQL然后就在同一个进程里查接收方在线状态再通过网关把消息推出去。链路短代码量也少看起来没毛病但真实跑起来就慢慢变味了。第一个亏是线程模型混乱。业务服务既要处理HTTP接口又要处理长连接的推送请求。一旦在线用户多起来转发时调网关接口被堵住整个业务服务的线程池都会被拖死。最典型的现象是IM接口变慢用户资料接口、会话列表接口跟着一起慢定位问题的时候根本分不清是业务问题还是推送问题。第二个亏是离线逻辑没法收敛。产品上了多端同步、离线消息、已读回执之后离线存储、消息序号、ACK确认这些逻辑越堆越多。每次改一个需求都要在这个混合体里小心翼翼地绕上线前比拆炸弹还紧张。第三个亏是没法独立治理。想针对“转发”做限流、熔断、灰度发布但这些代码和业务代码绑在一起动一处牵全身。要让转发服务单独扩容更是不可能。所以后来拆分“消息转发子服务”不是跟着微服务的风而是被现实问题逼的。1.2 拆分后的职责边界只做三件事拆分之后的链路变成了这样客户端 - 接入网关长连接 - 消息业务服务校验、持久化、生成序号 - MQ - 消息转发子服务查路由、查在线状态、投递/离线组织 - 接入网关下发 - 客户端我给转发子服务定的边界就三条不负责业务校验不负责用户资料不负责人际关系。它只需要根据消息体里的会话ID拿到会话成员列表再查每个成员当前挂在哪个网关节点、有哪些在线设备最后把消息推给在线设备把离线的部分写入离线队列。边界线看着朴素但它带来的直接好处是转发子服务下游依赖只有Redis里的路由表/在线状态以及推送网关这两类组件。业务再怎么变只要这两类数据不出问题这个服务基本不用动。2. 单聊、群聊、离线一条消息在转发子服务里的完整行走路线2.1 单聊在线优先、离线兜底的投递逻辑单聊流程上是最简单的。消息到达转发子服务后先查会话成员比如会话ID是s_10001里面只有A和B两个人。然后分别查A和B的在线设备表每台设备都记录了当前挂在哪个接入网关节点上。假设B有两台设备在线一台手机一台桌面端那就往两个网关节点各推一次如果某一台不在线就把消息写入B的离线消息ZSET。这里容易踩的一个坑是“设备在线状态”的过期时间。早期我把在线状态存RedisTTL设了90秒目的是让异常掉线的设备尽快清理掉。但在网络抖动时某台设备如果被误判离线转发服务会把它当作离线处理写进离线队列。等用户重新上线离线拉取和网关补推同时到达客户端就会看到消息重复。后来我不得不在客户端做消息ID去重服务端也把在线状态改成“心跳续租主动过期”的双机制而不是单纯依赖TTL。2.2 群聊写扩散还是读扩散群聊是我实践中最容易出问题的地方。核心问题是怎么把一个成员的消息分发给群里所有成员行业里有两种思路写扩散和读扩散。写扩散就是消息来一条为群内每个成员写一条待收记录。好处是每个成员拉消息时很快直接取自己的记录就行坏处是写放大非常严重。500人群里发一句话要写500份1000人群要写1000份群一多、消息一频繁写入量极其惊人。读扩散就是群里发消息只存一份群消息记录。每个成员的未读数通过偏移量来算打开会话时按偏移量拉取。好处是写入完全不放大坏处是拉取时必须实时计算查询压力大。我自己的落地做法是折中处理小于等于200人的群用写扩散按成员各写一份大于200人的群用读扩散只存群消息中心和成员游标。200这个阈值不是拍脑袋定的我做过分场景的写入与查询成本对比群人数写扩散单条消息总写次数读扩散单条消息写次数读扩散拉取开销50人50次1次非常低200人200次1次低500人500次1次中1000人1000次1次较高从压测数据看200人附近是两条策略综合响应时间的交叉点。小于200人时写扩散的存储代价完全可以接受而且成员拉消息体验好大于200人时写放大会让存储和推送扛不住必须转读扩散。2.3 离线消息与确认语义顺序比结果重要离线消息如果只按“在不在线”来切那边界情况会很多。我的做法是无论用户在线与否消息在业务服务阶段就已经分配好了一个会话级递增序号。转发时判断设备在线状态在线就走网关推送不在线就进离线ZSET以消息序号作为score。这样同一会话里所有消息都有明确的先后关系离线消息拉取时按序号批次拉取天然有序。然后是ACK确认。客户端收到消息后会回一个已收确认服务端根据客户端已经确认到的最大序号来判断离线消息可以清理到哪个位置。注意这里一定要用“连续确认序号”而不是“逐条确认”。因为网络环境里逐条确认会带来大量无用请求而且客户端收到的顺序未必完全一致用连续确认可以大幅减负。3. 消息顺序、路由表和MQ转发子服务的具体技术落地3.1 消息不直接走HTTP改走MQ的理由很多人做IM的第一步习惯把发送消息做成HTTP同步接口。但转发这块如果直接用同步调用业务服务的吞吐上限直接取决于转发服务的能力。而转发服务又需要频繁查询在线状态任何一次Redis抖动都可能把业务服务拖死。把MQ引进来之后业务服务只需要保证“消息可靠地写到MQ”就可以给客户端返回“已受理”。转发子服务从MQ拉消息按会话ID散列到固定worker再执行后续投递。MQ在这里的价值有三个削峰填谷让转发消费速度平滑可控异步解耦业务服务不再依赖转发服务的实时响应天然缓冲即使下游网关抖动消息也能先停在MQ里不会造成请求丢失。顺序性问题则是依靠MQ的分区机制来解决。Kafka的partition key、RocketMQ的顺序队列都是把同一会话的消息路由到同一个分区或队列。我的实现是选用RocketMQ因为它自带延迟消息和顺序消费的API对IM场景更友好。3.2 Redis里的路由表与会话序号设计路由表是转发子服务的核心数据。我用Redis来存键设计大致是这样// 会话成员列表会话维度 session:{sessionId}:members typeset // 会话基础信息 session:{sessionId} typehash (ownerId, type, maxMembers, msgSeq) // 用户当前在线设备与网关映射用户维度 user:{userId}:devices typehash deviceId - gatewayId // 设备心跳与最近活跃时间 device:{deviceId} typehash (userId, gatewayId, lastHeartbeat) // 用户离线消息队列 offline:{userId} typezset scoremsgSeq membermessageId注意两点第一会话成员列表用set因为IM会话成员不允许重复查成员时直接得到去重结果第二离线消息一定要用zset而不是list因为zset天然按序号排序拉取时可以用ZRANGEBYSCORE按区间取清理时按score删非常合适。这里还要配一个消息序号的生成策略。业务服务在真正写MQ之前先对每个会话做一次INCR操作// 业务服务 seq Redis.INCR(seq:single: sessionId) message { sessionId, sender, content, seq, ts } producer.send(im_msg_single, sessionId, message)转发子服务消费时就能复用这个seq来做排序。为什么不在转发阶段再生成序号因为消息可能在业务服务阶段被重试重试时如果序号重新生成后面消费的顺序就乱了。所以顺序号的分配点必须尽量靠前、尽量幂等。3.3 消费端如何使用Hash保证会话有序转发子服务消费MQ时我会在本地维护一个“会话到worker”的映射用会话ID的哈希值固定到某一个worker线程或协程。这样同一个会话的所有消息永远由同一个worker处理后再投递给网关顺序不会乱。核心消费伪代码大致如下// 转发子服务消费示例为Go描述 func onSingleMessage(message *Message) { // 根据会话ID取路由 members : redis.SMembers(session: message.SessionId :members) // 遍历成员查在线设备 for _, uid : range members { devices : redis.HGetAll(user: uid :devices) for deviceId, gatewayId : range devices { // 判断设备是否活跃 if isActive(deviceId) { pushToGateway(gatewayId, message) } else { redis.ZAdd(offline:uid, message.Seq, message.Id) } } } }这里有个容易忽略的问题如果同一个会话的某条消息处理失败我不会在线程里无限重试而是把消息放回一个本地重试队列延迟1秒、5秒、30秒分档重试。重试期间同一会话的其他消息怎么办我把它们先暂存在这个会话的“等待队列”中确保当前消息处理完之后才处理下一条避免重试时消息乱序。4. 高并发场景下的压测数据与群聊扇出优化4.1 第一轮压测暴露出的问题说多少理论都不如压测结果来得直接。提前说明以下数据来自我们的测试环境机器配置是8核16G的容器。测试场景从简到繁逐步加压。压测场景测试规模峰值吞吐P99延迟暴露问题单聊在线消息1万终端、1000并发发送8200 msg/s82ms无明显瓶颈单聊消息50%离线同上规模6100 msg/s150ms离线写入开始竞争500人群聊2000成员全员在线单群消息处理延迟P99 450ms消费滞后明显群聊扇出放大1000人群聊4000成员全员在线群消息吞吐上不去本地内存持续上涨必须拆批处理单聊场景的瓶颈不大因为一条消息最多投递给个位数的设备。但群聊一到500人以上问题就非常明显一条群消息会产生500条投递指令如果每条指令都携带完整消息体传输成本被放大了500倍。转发服务的CPU、网络带宽、内存立刻吃紧。4.2 群聊扇出优化的关键先聚合再扇出群聊优化的核心思想是不要让每条消息都变成N条重复的推送指令。我的落地做法是按网关节点做聚合。转发服务先把群内所有在线成员按“所在网关节点”分组把同一个网关下的所有成员合并到一个批量推送请求里消息体只传一份目标列表传设备ID数组。举个具体例子一个500人群里有300人分布在同一台网关节点上那么原来需要300条指令现在只需要给那台网关发1条批量指令。转发服务的CPU和网络开销立刻降下来。实测下来做聚合之后单节点处理群消息的能力直接翻了三倍P99延迟从450ms降到180ms左右。另一个优化点是把“离线写入”从投递链路里拆出去。既然消息是写扩散读扩散混合模式离线成员其实不需要立即逐个写离线队列可以先把“未读偏移量”更新到会话游标等成员下次上线时按游标去拉取。这个做法单独针对大群能极大降低离线部分的写入压力。4.3 背压、流控与降级的配置IM系统里发生消息堆积不可怕可怕的是堆积之后转发服务内存先爆掉。我在这块踩过坑所以后来特别强调“背压”设计。转发子服务消费MQ时本地会有一个待处理队列。队列长度不能是无上限增长的我设置了一个硬上限比如10万条缓冲超过后不再从MQ拉取新消息。宁可让消息滞留在MQ里也不能让进程内存失控。消费速率则由一个单独的信号控制如果Redis查询耗时升高消费线程会主动降低拉取频次给下游喘息时间。熔断也是必须的。当Redis连续失败达到阈值转发服务对“在线状态查询”这个动作直接降级先按“全部视为离线”处理把消息写入离线队列等Redis恢复后再重新在线投递。这样做会带来一定的延迟但避免了整个转发链路雪崩。这个降级策略的细节我在下一节的线上故障里会具体讲。5. 线上故障“failed to fetch dynamically”背后的完整排查链路5.1 故障现象成功率报警加一条前端日志那次故障发生在一个工作日下午。监控先报警消息发送成功率低于99%持续了大概七分钟。紧接着值班群里有人贴了客户端截图用户看到的是一个模糊的报错提示控制台里能看到一行failed to fetch dynamically看到这个报错前端同学第一反应是查构建产物、查CDN因为这条错误太像JavaScript动态导入分包失败的样子了。但查了一圈没有任何前端发布记录CDN也完全正常。5.2 排查推进把报错往上追我介入以后做的第一件事不是看业务日志而是看网关的链路日志。因为这个“动态获取”其实不是JS的import而是IM SDK在做“动态获取接入参数”时的报错。SDK里有个逻辑每次建立会话窗口时会调用一次配置接口从服务端动态拉取转发子服务的集群路由信息节点列表、协议版本、灰度标记。如果这次拉取失败SDK就会统一封装成这条错误文案。也就是说用户看到的“failed to fetch dynamically”其实是SDK在启动阶段没能从服务端拿到“动态接入配置”并不是聊天消息本身推不出去。定位到这里方向已经从前端转到了服务端。5.3 根因链条短TTL路由配置在热点回源下的连锁反应服务端排查时我看到了这样一串关键数据配置聚合服务的CPU升高但不是压倒性的转发子服务的“动态路由接口”内网专用于提供节点列表的接口成功率跌到了70%Redis里路由配置缓存Key的数量从正常时的1万个骤降到0。一步步往下追问题出在一个“短TTL加动态刷新”的路由配置机制上。我们的路由配置缓存设定是60秒TTL也就是说每60秒就要由配置聚合服务去转发子服务动态拉一次最新节点列表再写回Redis。这个机制本身是为了灰度发布时能快速生效但它的设计有个缺陷缓存过期后的首次请求会触发回源如果同一瞬间有大量请求同时触发回源转发子服务就被打满了。故障的触发点是一位同事发了一封全员通知导致Web端大量长连接几乎同时重连。重连之后每个SDK都会立刻创建会话窗口并发起动态接入配置的拉取请求。这一波流量是多少是平时的80倍。配置聚合服务的本地缓存刚好在那一刻过期于是全部请求穿透到转发子服务的动态路由接口。更致命的是转发子服务的动态路由接口是同步查Redis再拼装集群最新状态的。而Redis在那一段时间正好因为热点Key击穿在处理大量重试响应出现批量超时。接口的线程池被打满后续动态获取全部失败。SDK拿不到配置就直接报出了那行让人误会的错误。5.4 修复方案动态获取的兜底与降级这个故障从表面看是“动态获取失败”根子上却是三件事回源没有保护、缓存只有短TTL没有兜底、服务端和客户端都缺少降级。修复方案分成三层服务端动态路由接口增加“兜底静态配置 异步重建”。本地缓存永远保留上一份可用快照即使动态刷新失败也返回上一份可用数据不抛错除非本地快照都不存在才返回失败。这样做保证动态失败不会引发大面积故障。配置聚合服务增加“缓存续期”逻辑如果Redis不可读不得把失败透传成回源而是给本地缓存续期90秒让转发子服务缓一口气。客户端SDK对“动态获取失败”增加重试指数退避5秒、30秒两个档次最多重试三次。不要第一次失败就直接把错误抛给用户。这次故障之后我也把“动态配置”这一类接口统一加了一个规则凡是配置类数据必须有静态兜底绝不允许因动态源故障而影响核心链路。配置可以不是最新的但不能拿不到。6. 对消息转发子服务的几点个人总结6.1 “消息发送成功”的定义最好写代码前就想清楚在IM项目里客户端提示“消息已发送”和服务端真正把消息推到对方手里中间差着十万八千里。我踩过的教训是不要把“消息写进MQ”当作成功也不要把“消息推给网关”当作成功你要根据产品的需求去定义清楚是“服务端已存储”就算成功还是“对端已收到ACK”才算成功。这个定义直接影响转发子服务的状态机设计。6.2 转发子服务可以不承载业务但必须承载足够的可观测性转发子服务是整个IM消息链路的腰部任何网关抖动、Redis抖动、路由表异常都会在这里汇聚。所以日志一定要带完整的会话ID、消息序号、设备ID、网关ID并且要和客户端的请求链路ID打通。这次“dynamic fetch”故障能快速定位靠的就是SDK里保留了完整的链路上下文。否则拿着一个报错文案去做前端排查做一天也找不到根因。6.3 架构设计时留下的坑压测时往往发现不了短TTL缓存的设计在平时压测完全正常因为压测不会模拟几千个Web端同时重连的场景。但真实世界里全员通知、热点事件、客户端版本强制升级都可能瞬间制造出远超压测的流量规模。所以对待配置类数据我现在的原则是能静态不全动态能缓存就缓存动态刷新失败必须降级绝不能把“动态获取”做成一个硬依赖。这个原则适用于IM也适用于任何有动态路由、动态配置的服务。做消息转发子服务这两年我对这个模块最大的感受就是它不负责创造业务价值但所有业务价值都要靠它来送达。把这个“送信人”养稳整个IM项目就稳了一大半。
返回列表