从零自研Java原生IM系统:长连接、消息可靠性与高并发实践
2026/9/15 10:02:26 网站建设 项目流程

先说一个很多人容易忽略的事实:IM 即时通讯系统,表面上是个“聊天工具”,本质上是一套“长连接管理 + 消息可靠性 + 高并发推送”的组合体。网上一搜 Java IM 源码,出来的开源项目少说也有十几个,可真拿到公司内部用,面对客服分流、组织架构同步、消息审计、多端登录这些定制需求时,你会发现改别人的代码,比从零写一套还难受。

我去年帮一家创业公司从零落地过一套基于 Java 原生技术栈的 IM 系统,这里说的“原生”,不是让你从 Socket 手写开始,而是以 JDK 为底座、以 Netty 为网络基础设施,自己掌控协议、会话、存储、推送、部署这整条链路。这套方案的好处是:每一层你都看得懂、改得动,出问题能顺着代码一路查到根因,而不是在黑盒里猜。这篇文章会把通信链路设计、消息可靠性机制、群聊分发策略、部署演进方案,以及我实测踩过的几个大坑完整拆开讲,适合正准备自研 IM、或者想彻底搞懂 IM 底层原理的 Java 后端开发者参考。

1. 为什么还要坚持 Java 原生,而不是直接套开源框架?

1.1 开源 IM 的“拿来主义”陷阱

先聊一个很现实的问题:OpenIM、悟空 IM、芋道 IM 这些开源项目,功能截图一个比一个漂亮,代码量也是动辄几十万行起步。很多人拉下来编译通过后,觉得“已经到手了”,结果一接业务就发现完全不是那么回事。

IM 系统不是 CRUD 管理系统,它的复杂度集中在连接状态、消息顺序、离线补偿、多端同步这些“实时性”细节里。当你需要改一个群公告的推送逻辑,或者调整消息撤回的时序策略时,必须把别人的连接管理、会话路由、存储模型全部吃透才能动手。我见过最典型的例子:某个开源项目把消息表设计成单表不分区,在线用户一到两万,消息表直接卡死 MySQL,而你想改它的存储层,得先理解它散落在十几个类里的 DAO 调用链——这种改造工作量,真的不如自己重写。

1.2 什么才算真正的“Java 原生”实现

我讲的“原生”,落点其实是这条技术栈组合:

  • JDK 自带 NIO 与并发工具,负责底层 IO 模型与线程调度;
  • Netty 只作为高性能网络通信库,提供 Reactor 线程模型、编解码框架、拆包粘包处理;
  • 业务链路(协议定义、会话管理、消息存储、推送策略、接入层路由)全部由自己实现。

很多人纠结“用了 Netty 还算不算原生”,我觉得这是把框架和基础设施搞混了。Netty 的定位是网络编程库,不是 IM 框架,它没有替你决定会话怎么存、消息怎么排序、离线消息怎么拉取。真正决定 IM 系统灵魂的,恰恰是你自己写的这一层业务链路。这也是为什么我推荐自研 IM 时,选 Netty 而不是选一个完整的 IM 开源框架——给你一块地基,和给你一栋不能拆墙的房子,是完全不同的两码事。

1.3 哪些场景适合自研,哪些场景别碰自研

这里给一份我自己的选型判断表,方便你对号入座:

场景特征推荐方案原因
内部 OA/IM,定制化需求极多自研(Java 原生 + Netty)可控性最强,改造成本最低
只有 2~4 周交付时间,后端人手不足直接接云厂商 IM 或开源项目自研连基本链路都跑不完
面向 C 端,日活百万以上自研,但要做好架构分层开源方案很难支撑这种规模与定制
业务是标准聊天,无特殊合规要求开源项目二次开发省时省力,前提是别碰核心链路

我自己踩过一次教训:有一个项目本来只需要做 App 内嵌客服聊天,我按“标准 IM”的规模设计了群聊、已读回执、多端同步,结果团队三个人干了三个月,客服那边只用到了单聊 + 离线消息 + 工单关联。后来我把这套思路沉淀成一句话:先搞清楚你的 IM 是“聊天工具”还是“业务系统里的一条实时通道”,这决定了你是要造一台车,还是只装一个轮子。

2. 通信链路设计:长连接、协议拆包与心跳保活

2.1 通道选型:TCP 长连接、WebSocket 还是 HTTP 轮询

IM 的实时性全靠连接通道撑着,通道选型直接决定上下行延迟和服务器资源消耗。我在实际项目里遇到过三种通道混用的情况,这里直接给结论:

  • 移动端 App、桌面客户端、服务端之间:用 TCP 长连接。Netty 维护起来最顺手,控制力最强,适合自定义私有协议。
  • 浏览器端、H5 页面:用 WebSocket。浏览器只认这个,TCP 裸连在网页端根本推不进去。
  • 低实时性场景(比如通知类消息、邮件提醒):HTTP 短轮询或者 Server-Sent Events(SSE)足够,不要浪费长连接资源。

通道这块最容易犯的错误,是把 WebSocket 当成长连接的全部。其实 WebSocket 的底层也是 TCP,它只是给 TCP 加了一层浏览器友好的握手与帧封装。如果你在服务端统一收敛成 TCP 长连接,网关层做协议转换(TCP 协议与 WebSocket 协议互转),上层业务逻辑就不用关心客户端到底是从 App 来的还是从浏览器来的。

2.2 自定义私有协议:消息头到底该放哪些字段

既然走 TCP 长连接,就必须设计私有协议。很多新手喜欢直接用 JSON 字符串加换行符当协议,开发期确实方便,一上生产就出问题:JSON 解析耗 CPU、消息体内含换行符会拆包错乱、没法做高效的二进制扩展。

我们最终用的协议方案是“二进制消息头 + 可序列化消息体”,消息头用固定 24 字节:

+--------+--------+--------+--------+----------------+----------------+----------------+ | magic | version| command| flags | sequenceId | bodyLength | reserved | | 2 bytes| 1 byte | 1 byte | 1 byte | 8 bytes | 4 bytes | 7 bytes | +--------+--------+--------+--------+----------------+----------------+----------------+
  • magic:魔数,固定 0xAC 0xED,用来快速识别非法连接,防止非 IM 客户端乱发数据打爆服务端。
  • command:命令字,比如 0x01 登录、0x02 心跳、0x03 单聊消息、0x04 群聊消息。
  • flags:标志位,比如是否压缩、是否需要 ACK、是否是离线消息补偿。
  • sequenceId:客户端生成的递增序列号,用于消息去重和服务端响应关联。
  • bodyLength:消息体长度,这是拆包的关键字段。

消息体默认使用 JSON 序列化,因为业务组维护成本低,后续如果遇到性能瓶颈,再针对高频消息(比如文本聊天)改成 Protobuf。这里面的经验是:协议头要固定、要紧凑,消息体要灵活、要好维护,不要为了追求极致的性能把体量很小的 JSON 也换成二进制,工程上的收益不值当。

2.3 粘包拆包:Netty 的 LengthFieldBasedFrameDecoder 参数怎么调

TCP 是流式协议,没有消息边界。客户端连续发送多条消息时,服务端可能一次读到半条消息,或者一次读到好几条粘在一起。Netty 里解决这个问题最优雅的方式就是 LengthFieldBasedFrameDecoder。

我们在代码里是这么配置的:

new LengthFieldBasedFrameDecoder( 1024 * 1024, // maxFrameLength:单条消息最大 1MB,超过直接报错 0, // lengthFieldOffset:长度字段从第 0 字节开始 4, // lengthFieldLength:长度字段占 4 字节 2, // lengthAdjustment:长度字段之后还有 2 字节的 command 0 // initialBytesToStrip:不剥离任何字节,让后续 handler 拿到完整协议头 )

有个细节必须强调:lengthAdjustment 很多人搞不明白。如果协议设计是 [4 字节长度][2 字节 command][N 字节 body],长度字段的值是 N,那么实际帧的总长度 = 4 + 2 + N。Netty 拿到长度字段后,要加上 lengthAdjustment 才能算出完整帧的结束位置。这里的取值是 2,对应当前帧“长度字段之后、body 之前”的剩余字节数。

调完拆包器之后,不要马上联调,先用十六进制报文工具测一遍边界情况:半包、粘包、空 body、body 长度字段被恶意写大。我见过太多线上 IM 事故,根子都出在拆包边界没测透。

2.4 心跳保活:为什么客户端心跳间隔要比服务端探测间隔短

TCP 长连接不会自己一直活着,中间任何一层网络设备(路由器、防火墙、运营商 NAT)都可能静默掐断空闲连接。所以 IM 系统必须有心跳机制。

我们用 Netty 的 IdleStateHandler 实现:

// 服务端 ch.pipeline().addLast(new IdleStateHandler(60, 0, 0, TimeUnit.SECONDS)); // 60 秒没读到客户端数据,就触发 userEventTriggered // 客户端 ch.pipeline().addLast(new IdleStateHandler(0, 25, 0, TimeUnit.SECONDS)); // 25 秒没向服务端写出数据,就主动发一次心跳包

这里有个非常关键的经验:客户端的写空闲时间(25 秒)必须小于服务端的读空闲时间(60 秒)。原因很简单,客户端主动心跳是维持连接的“保鲜动作”,服务端被动探测是“兜底清理动作”。如果两者反了,客户端 60 秒才发一次,服务端 60 秒没收到数据就判定超时,很容易因为一次网络抖动导致连接被误杀,客户端还来不及补救,体验就会断崖式下降。

心跳包本身要做得越轻量越好,服务端收到后只需要更新最近活跃时间,不用落库、不用回业务消息,返回一个 4 字节的 ACK 即可。如果心跳要携带业务数据,就说明你的设计已经跑偏了。

2.5 断线重连:指数退避必须加随机抖动

客户端断线重连,最忌讳的是所有客户端在同一时刻无限重试。想象一个场景:晚上 12 点 Nginx 超时断开了一半连接,第二天早上 8 点用户集中打开 App,重连请求瞬间把服务端打崩——这不是网络问题,这是重连策略问题。

我用的方案是“指数退避 + 随机抖动”:

retryDelay = Math.min(60, baseDelay * (2 ^ retryCount)) + RandomUtil.randomInt(0, 5000); // baseDelay 初始 1 秒,retryCount 最大重试次数

从 1 秒开始,第二次 2 秒,第三次 4 秒,依次翻倍,上限 60 秒,每次再加上 0~5 秒的随机值。这样既能保证大部分客户端在网络恢复后 1 分钟内连回来,又不会形成重连请求的流量尖峰。

还有一个容易漏的细节:客户端重连成功后,必须做一次增量同步,把断线期间漏掉的消息拉回来。否则重连只是“连上了”,消息还是对不齐。这个增量同步的位点,就是消息可靠性设计里的 seq 概念,下面展开聊。

3. 消息可靠性:链路随时会断,但消息一条都不能丢

3.1 消息生命周期:一条消息从发出到已读的完整链路

客户端小李给小王发一条文本消息,这条消息会经历这样一串流程:

  1. 小李的客户端生成全局唯一 messageId(UUID 或雪花 ID),把消息发送到服务端接入层;
  2. 接入层校验完整性后,把消息投递到消息处理模块,落库存储;
  3. 存储成功后,服务端给小李回一条 ACK:消息已收到;
  4. 服务端查询小王的在线状态,如果在线,推送到小王所在网关节点,小王客户端收到后回执 ACK;
  5. 如果小王离线,消息进入离线消息表,等小王下次登录后按 seq 增量拉取;
  6. 小王看到消息后,客户端上报已读回执,单聊场景服务端更新会话已读位点。

这条链路每一环都可能断。可靠性设计的核心就是在每一环都提供“重试 + 幂等 + 对账”的能力。

3.2 ACK 机制:发送方、接收方、服务端三方各自的职责

ACK 不是只有一种。我一开始设计的时候只做了“客户端到服务端”的 ACK,结果消息发送方永远不知道消息到底送达没有。后来拆成三类:

  • 上行 ACK:客户端发消息,服务端持久化成功后回的业务确认。客户端超时没收到就重发,重发带同一个 messageId。
  • 下行 ACK:服务端推消息给接收方,接收方客户端收到后回的确认。服务端如果没收到这个 ACK,会重推几次,超过次数就转入离线消息。
  • 已读回执:接收方真正看到消息后上报,这个改变的是会话的已读位点,不影响消息是否投递成功。

这里有个性能优化点:下行 ACK 没必要逐条回,客户端可以把多条消息的 ACK 合并成一个批量 ACK 包,服务端用一个数组接收。否则群聊场景下,几十条消息刷过来,ACK 风暴能把服务端小包打满。

3.3 消息序号(seq)与离线消息拉取

每个用户维护一个单调递增的消息序号 seq,这个 seq 由服务端统一分配。为什么不能由客户端自己生成?因为多端登录的时候,手机和电脑同时发消息,客户端本地序号会冲突,排序会乱。服务端分配 seq 的核心作用有三个:

  • 消息排序依据;
  • 离线消息增量拉取的游标;
  • 重复消息去重的参考位点。

离线消息的存储,我用的是一张独立的离线消息表:

CREATE TABLE `offline_message` ( `id` bigint NOT NULL AUTO_INCREMENT, `user_id` bigint NOT NULL COMMENT '接收方用户ID', `msg_seq` bigint NOT NULL COMMENT '该用户的全局消息序号', `msg_id` varchar(64) NOT NULL COMMENT '消息全局唯一ID', `sender_id` bigint NOT NULL DEFAULT '0', `content` text COMMENT '消息内容', `create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (`id`), KEY `idx_user_seq` (`user_id`, `msg_seq`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

客户端登录成功后,把本地已确认的最大 seq 上报,服务端查出这之后的所有离线消息下发。这里有个大坑:离线消息表会无限膨胀,必须在每次拉取成功后清理过期数据。我们的策略是“拉取确认后延迟 5 分钟删除”,防止客户端拉取完但还没落本地时崩溃导致数据丢失。

3.4 消息幂等:重发可以重来,但不能重复入库

ACK 超时重发机制上线后,最直接的问题就是一条消息可能被发送方重复投递多次。如果没有幂等设计,接收方就会看到两条一模一样的消息。

幂等的关键就是 messageId。服务端收到一条上行消息时,先查 Redis 里有没有这个 messageId 的处理记录:

Boolean first = redisTemplate .opsForValue() .setIfAbsent("im:msg:" + messageId, "1", Duration.ofHours(24)); if (first == null || !first) { // 重复消息,直接返回 ACK,不重复落库 return ack(messageId); }

Redis 的 SETNX 天然适合做这种一次性去重,但要注意给 key 设置过期时间,避免消息去重表无限膨胀。如果消息量极大,还可以把去重从 Redis 挪到本地 Caffeine 缓存 + 数据库唯一索引双层保证。我个人建议消息表对 messageId 建唯一索引,这是最后一道兜底,Redis 挂了也不会重复入库。

4. 群聊与高并发场景:别把单机方案的思路硬套到分布式

4.1 群聊消息的两种分发策略:扩散写与读扩散

群聊是 IM 系统里最容易出性能问题的场景。一个 1000 人的群,有人发一条消息,是往群里每个人的收件箱里都存一份,还是只存一条、大家各自拉取?这两种方案各有代价:

  • 扩散写:消息到达服务端后,给群内每个成员的消息表都插一条记录。好处是接收方离线消息拉取逻辑简单;坏处是写放大严重,1000 人群就是 1000 次写。
  • 读扩散:消息只存一份,接收方拉取时动态合并自己加入的所有群的时间线。好处是写次数少;坏处是拉取逻辑复杂,需要合并群消息和单聊消息,分页很容易乱。

我们实盘的经验是混合策略:群成员 200 人以下用扩散写,因为写量可控,接收方体验最好;200 人以上的大群用读扩散,避免频繁的大规模写放大。这个阈值不是拍脑袋定的,我们用压测数据算过,MySQL 单机在混合读写场景下,200 人左右的扩散写还能保持在 SLA 之内,再往上就会明显拖慢消息时延。

4.2 会话路由:用户在不同网关节点之间,消息怎么找到对方

单机版 IM 不需要考虑路由,所有连接都在一个 JVM 里,直接内存寻址。一旦接入层部署了多个 Netty 节点,用户 A 连在 node1,用户 B 连在 node2,A 给 B 发消息时,node1 怎么知道该转发给 node2?

我们的方案是用 Redis 维护一张在线路由表:

key: im:route:{userId} value: { nodeId: "node-1", channelId: "xxx", lastHeartbeat: 1699999999 } TTL: 90 秒(依赖心跳续期)

消息发送流程变成:

  1. 客户端 A 发消息到 node1;
  2. node1 解析接收方 userId,查 Redis 路由表;
  3. 如果接收方在线且在同一节点,直接 channel.writeAndFlush;
  4. 如果接收方在别的节点,走内部 RPC 转发到目标节点;
  5. 如果接收方离线,进入离线消息流程。

跨节点转发的内部链路,小规模用 Netty 节点间的内部长连接通道就够了,也可以用 RocketMQ / Kafka 做异步解耦。我建议接入层节点少于 10 个的时候,直接用内部 RPC,少引入一个 MQ 中间件,就少一类故障。

4.3 在线状态管理:上下线、踢人、多端登录的底层逻辑

在线状态不是一个布尔值,而是“用户与网关节点”的映射关系。用户上线时,客户端通过接入层鉴权成功后,服务端往 Redis 写路由信息;用户下线时,客户端主动发下线请求,或者服务端探测到连接断开后,把路由信息标记下线。

多端登录是另一个容易踩坑的点。同一个账号在手机、PC、网页同时在线,在线状态必须从“用户维度”细化到“设备维度”。我们的方案是允许同账号最多 3 个端同时在线,每一端有独立的 channelId 和 deviceType。广播消息时,服务端先查用户的所有在线设备,逐设备下发;如果某个设备的连接已经失效,只移除该设备的路由,不影响其他端。

这里要特别提醒:千万不要把在线状态设计成简单的 Redis String 覆盖写。A 端上线把 userId 对应的 value 覆盖成 A 的 channel,B 端一上线把 A 顶掉,用户开个网页版 IM 就把 App 端挤下线,绝对会被业务方骂死。

5. 部署落地方案:从单机演示到可灰度上线的演进路径

5.1 单机最小闭环:一套 Docker Compose 跑起来

刚开始做功能联调或给老板演示的时候,不需要一上来就搞 K8s,一个 docker-compose.yml 就能把整套环境拉起来。最小闭环包含三个容器:IM Server(Spring Boot + Netty)、MySQL、Redis。

version: "3.8" services: im-mysql: image: mysql:8.0 container_name: im-mysql environment: MYSQL_ROOT_PASSWORD: root123 MYSQL_DATABASE: im_db ports: - "3306:3306" volumes: - ./sql:/docker-entrypoint-initdb.d command: --character-set-server=utf8mb4 --collation-server=utf8mb4_unicode_ci im-redis: image: redis:7.0 container_name: im-redis ports: - "6379:6379" im-server: image: openjdk:17-jdk-slim container_name: im-server depends_on: - im-mysql - im-redis volumes: - ./target/im-server.jar:/app/im-server.jar ports: - "8080:8080" - "9000:9000" environment: SPRING_DATASOURCE_URL: jdbc:mysql://im-mysql:3306/im_db?useUnicode=true&characterEncoding=utf8 SPRING_DATA_REDIS_HOST: im-redis entrypoint: ["java", "-Xms512m", "-Xmx512m", "-jar", "/app/im-server.jar"]

这里有一个很实用的建议:把建表 SQL 放在 ./sql 目录,挂载到容器的 docker-entrypoint-initdb.d,MySQL 容器首次启动时会自动执行,省得每次部署都手动导表。IM Server 暴露两个端口,8080 走 HTTP 用于健康检查和管理接口,9000 是 Netty 的 TCP 长连接端口。

5.2 多节点部署:从单机到网关层无状态化

当在线用户数上来之后,单机 Netty 的连接数会先触顶。Linux 单进程可以维持的连接数理论上很高,但受限于文件描述符、线程调度、GC 压力,实际生产中单节点支撑 5 万~10 万长连接已经是比较稳的上限了。

多节点部署时,接入层要做到“无状态”:

  • 连接状态全部落在 Redis,节点本身不保存用户路由;
  • 客户端启动时通过 Nginx / SLB 做负载均衡,连接到任意可用节点;
  • 节点重启前做优雅下线,客户端感知后自动重连到其他节点。

架构上会变成这样一条链路:

客户端 App / Web │ ▼ 接入层 Nginx / SLB(四层负载均衡,转发 TCP 长连接) │ ├──► IM Server Node 1 (Netty) ├──► IM Server Node 2 (Netty) └──► IM Server Node 3 (Netty) │ ▼ Redis(路由表 + 消息去重 + 在线状态) │ ▼ MySQL(消息存储 + 离线消息 + 用户关系)

这里有一个关键配置:如果用的是 Nginx 四层转发 TCP 长连接,必须把 proxy_timeout 调大,默认 60 秒的空闲超时会把你 IM 的心跳连接全部切断。我见过一次线上事故,IM 消息经常发不出去,排查到最后发现是 Nginx 的 stream 模块默认空闲超时把长连接全断了。

5.3 优雅停机:重启服务端时,怎么做到用户无感知

IM 服务端发布重启最怕什么?怕存量连接被直接 kill,客户端上一秒还显示在线,下一秒全部断线重连。如果几十万客户端同时重连,重新接入的请求洪峰能瞬间打满 CPU。

优雅停机要分三步走:

  1. 从负载均衡摘掉节点流量,不再接收新连接;
  2. 对存量连接广播“服务端即将维护”的推送消息,同时标记当前节点状态为 draining;
  3. 等待存量连接处理完当前请求,或者等待一个最大宽限期(比如 30 秒),然后关闭 EventLoopGroup。

Spring Boot 里可以用 ApplicationListener 监听 ContextClosedEvent,配合 @PreDestroy 做资源回收。Netty 侧的关闭顺序也很有讲究:先 channelGroup.close(),再 workerGroup.shutdownGracefully(),最后 bossGroup.shutdownGracefully()。顺序反了会出现新连接还没处理完就被 shutdown 的情况。

客户端那边也要配合:收到服务端的维护通知后,不要立即重连,等待几秒后再随机延迟重连,避免重启后再次形成连接风暴。

6. 实战中踩过的坑与性能调优笔记

6.1 把业务操作直接写在 EventLoop 线程上,这是最大的性能杀手

Netty 的 IO 线程(EventLoop)数量默认是 CPU 核数的两倍,它负责处理 Channel 的读写事件。很多人写第一个 Netty Handler 时,习惯直接在 channelRead 里查数据库、调外部接口,看起来没问题,实际上是把所有连接的 IO 处理全部堵在几个线程上。

正确的做法是:Handler 里只做编解码和消息路由,耗时操作丢给独立的业务线程池。我们的线程池模型是:

ExecutorService bizExecutor = new ThreadPoolExecutor( 16, 64, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(10000), new ThreadFactoryBuilder().setNameFormat("im-biz-%d").build(), new ThreadPoolExecutor.CallerRunsPolicy() );

CallerRunsPolicy 很关键:线程池满了之后,任务会回退给调用方线程(也就是 EventLoop)执行。这个策略看起来会导致 IO 线程被阻塞,但它能天然形成背压,防止任务无限堆积导致内存溢出。真实环境里,宁可短暂阻塞 IO,也不能让任务队列无上限增长。

6.2 ByteBuf 引用计数泄漏:堆外内存被吃光的隐形杀手

Netty 的 ByteBuf 使用堆外内存时,需要手动 release 释放引用计数。如果 handler 处理完消息后没有调用 ReferenceCountUtil.release,或者你用了 ctx.write 但忘了一端 release,堆外内存就会一点点泄漏,最终触发 OOM,而且崩溃前毫无征兆,只有 GC 日志里能看到堆外内存持续增长。

我自己排查这类问题的经验是:先打开 Netty 的泄漏检测级别:

System.setProperty("io.netty.leakDetection.level", "PARANOID");

这个参数会大幅降低吞吐量,不适合生产环境长时间开启,但定位泄漏源非常好用。日志里一旦出现 LEAK: ByteBuf.release() was not called before it's garbage-collected,跟着堆栈就能找到没正确的 handler。定位后改代码,再把检测级别调回 SIMPLE 或关闭。

6.3 消息乱序:同一个会话的消息不能从多个线程同时发出

分布式环境下消息乱序,是我后期才遇到的坑。同一个用户的消息,如果被分发到不同的 MQ partition,或者 Redis publish 的 channel 消费端并发处理,就会出现“消息 B 先到、消息 A 后到”的乱序。

解决办法是保证同一个用户的消息走同一个队列分片

  • RocketMQ / Kafka 场景:消息 key 按接收方 userId hash,保证同一 userId 的消息进同一个 partition;
  • 线程池场景:对目标 userId 取模,固定路由到同一个业务线程,从源头消除并发;
  • 数据库场景:消息落库后,推送动作不要并行,按 seq 串行发送。

压测时还要注意一个细节:TCP 层本身保证有序,但客户端如果同时维护了多条 TCP 连接(比如长连接断了又重连),新旧连接之间就可能并发收包,客户端必须按 seq 做一次重新排序。

6.4 数据库连接池被慢查询拖死

消息落库是 IM 系统里读写频率最高的操作。之前出现过一次故障:某天消息量突增,MySQL 慢查询日志里全是消息表的大分页查询,连接池被打满,线上消息发送集体超时。

排查后发现根源是分页拉取历史消息用了 LIMIT offset 的深分页写法,偏移量大了之后 MySQL 要扫描大量无用行。优化方案是改成基于 seq 的游标分页:

-- 不推荐 SELECT * FROM message WHERE session_id = ? ORDER BY msg_seq DESC LIMIT 20, 20; -- 推荐:带上最后一次拉取的 msg_seq SELECT * FROM message WHERE session_id = ? AND msg_seq < ? ORDER BY msg_seq DESC LIMIT 20;

另外,消息表的写入一定要做成批量插入,不要一条一条 insert。IM Server 里把 1 秒内到达的消息攒一把,一次 INSERT 带几十条 VALUES,配合 rewriteBatchedStatements=true,写入性能能提升一个数量级。

6.5 连接风暴不只是客户端的问题,服务端也要做自我保护

前面讲了客户端要指数退避,服务端也不能干等着被冲垮。我们的做法是接入层加一个简单的“半连接保护”:当当前活跃连接数超过最大阈值的 80% 时,对新连接直接返回“服务繁忙请稍后重试”,客户端收到这个响应后延迟较长时间再重连;同时监控连接建立速率,超过每秒预估值的 N 倍就触发告警,提醒运维扩容。

这种保护本质上不是提升吞吐,而是给系统留出缓冲时间。IM 系统的雪崩往往不是单点被打垮,而是连接风暴引发 CPU 被打满,CPU 打满导致心跳延迟,心跳延迟导致客户端误判断线,断线后再次触发重连,恶性循环。加一层自我保护,就是打断这个循环的第一道闸。

7. 一些压箱底的建议

整套系统落地之后回头看,最有价值的不是某个技术方案多精巧,而是“怎么把它组织成一条你能 hold 住的链路”。如果你也准备从零写一套 Java 原生 IM,我的建议是:第一版只做单聊 + 离线消息 + 在线状态这三件事,把这条主链路彻底跑通、压稳,群聊、已读回执、多端同步都是在这张网上挂能力。消息 seq 的分配和路由表的设计,一定要在最开始就统一收敛到服务端,不要允许客户端参与任何消息排序和去重逻辑,不然后面每次加功能都会在这里踩坑。

另外一个容易被忽略的点是监控。长连接系统和传统 HTTP 接口不一样,接口挂了能被调用方立刻感知,连接挂了可能没人知道。上线前至少要把这三个指标监控起来:活跃连接数、消息上行 TPS、消息端到端推送时延。哪天消息时延从 50ms 涨到 500ms,先看 GC 曲线,再看 MySQL 慢查询,最后看是否出现了连接泄漏,按这个顺序排查,绝大多数问题都能在用户投诉之前发现。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询