接手AI培训系统实时通讯模块重构之前,我一直觉得实时通讯这个方向没什么好聊的:无非是WebSocket、心跳、断线重连,网上教程一抓一大把。等真正把AI助教、课堂互动、批改结果流式下发、在线考试防作弊这些业务链路全部接进来之后,我才发现单靠一套WebSocket根本撑不住整个场子,MQTT的引入不是锦上添花,而是架构演进的必然结果。这篇文章把我这次技术选型、消息模型设计、性能调优和线上排障的完整过程整理出来,给正在做在线教育、AI培训类系统实时通讯模块的同学做个参考,尤其是那些在“WebSocket够用吗”“为什么要混用MQTT”“连接数上去之后怎么优化”之间反复纠结的人,这篇文章应该能帮你省掉不少弯路。
1. 先搞清楚业务再选协议:AI培训实时通讯到底在传输什么
很多人在技术选型时有个通病,先定协议再想业务,最后发现协议和场景八字不合。我这次把AI培训系统的实时通讯模块拆开来看,业务远不止一个“在线聊天”那么简单。
1.1 从“在线课堂”到“AI伴学”:实时链路不止聊天一条
传统在线培训系统的实时通讯,核心就是直播间弹幕、举手提问、课件翻页同步,这些场景用WebSocket确实足够。但我们这套AI培训系统最大的差异在于引入了大量AI实时能力,直接改变了通讯模型。
举几个真实场景:
- AI助教实时答疑。学员在练习过程中点击AI助教,后端大模型生成回答不是一次性返回的,而是流式吐字。这个“流式”必须通过实时通道推给前端,如果走HTTP轮询,要么延迟高,要么把后端打爆。
- AI作业批改与评分。学员提交一道主观题,AI批改引擎在3到5秒内完成分析,结果需要主动推送到对应学员端,同时还要通知教师的监听端。这不是简单的点对点通讯,而是“一条结果分发到多个订阅者”。
- 课堂实时数据看板。讲师端需要实时看到全班学员的答题进度、正确率分布、AI推荐的讲评题目,这些数据由多个内部服务产生,必须聚合后广播给讲师端。
- 考试防作弊监控。系统需要实时接收学员切屏、异常行为事件,同时推送监控指令给监考端,对弱网环境下的可靠性要求非常高。
这些场景共同点是:传输的数据不只是“人说的话”,更多是“系统事件”和“AI生成内容”。人说的话丢了可以重发,但系统事件丢一条可能导致看板数据对不上、批改结果缺失、AI对话中断,这就对通讯模块的可靠性、有序性、广播能力提出了完全不同的要求。
1.2 架构演进过程中的分水岭:为什么不能只用WebSocket
第一版我们确实只用了一套自研WebSocket服务,业务方把课堂互动、AI流式输出、系统通知全部往里塞。运行了半年,问题逐渐暴露:
- WebSocket是“连接型”协议,一端一个连接,天然适合“这个人到那个人”的交互,但当需要把一条AI批改结果同时推送给学员端、讲师端、数据大屏、消息中心时,应用层要做大量“订阅关系”管理,写着写着就变成了在业务里手搓消息中间件。
- AI培训系统内部有大量微服务在产生事件,Java后端、Python算法服务、Node.js辅助服务都要往实时通道里发消息。WebSocket的接入端是浏览器,服务端之间想通过它互相通信非常别扭,最后变成所有服务都直连一套长连接服务,耦合严重。
- 弱网场景下,WebSocket的断线重连、消息补发、离线消息处理,全部要自己实现,代码量不小,而且很难覆盖移动端和桌面端的各种网络切换。
到了这个节点,我意识到缺一个“消息中间件”的角色。实时通讯的选型不是“二选一”,而是分层:端到端的实时交互用WebSocket,系统内部的消息分发与事件路由交给MQTT,两者叠加才形成完整的通讯能力。
1.3 我的选型判断标准:连接模型、消息模型与可靠性
聊到选型,很多文章会列一堆对比表格,什么传输层、头部大小、协议开销,这些网上都能查到。我在实际决策中只看三个维度。
连接模型:WebSocket是长连接、双向、全双工,一条连接就是一条“管道”;MQTT是发布订阅模型,客户端连接Broker后,通过Topic订阅来收消息,发消息的人不需要知道谁在收。AI培训系统的实时场景,很大一部分是“事件产生方”和“事件消费方”解耦,MQTT天然契合。而学员与AI助教之间的会话交互,需要保持一个稳定双向连接,WebSocket的“直接性”更合适。
消息模型:如果业务是“A发给B”,用WebSocket;如果业务是“C产生了一条数据,D/E/F都要收到”,用MQTT的Topic做广播。AI批改结果推送给多端、课堂事件同步给看板、系统公告群发给所有在线学员,这些全是典型的“一对多”发布订阅,用MQTT一条代码都不用写,自带实现。
可靠性:WebSocket本身没有消息确认机制,应用层需要自己定义ACK和重传;MQTT有QoS 0/1/2三级,至少能保证“消息到达Broker”。对AI批改这类不能丢的消息,QoS 1就是现成保障。
提示:不要陷入“哪个协议更高级”的争论。真实架构里两者是互补关系,WebSocket负责终端交互,MQTT负责服务端消息枢纽,中间通过一个适配层打通。这样业务代码清晰,线上排障也容易定位问题在哪一层。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 消息模型与Topic设计:搭起WebSocket与MQTT的协作骨架
实时通讯模块最怕“只通不邮”,通道建好了,消息却没有统一规范,业务各写各的,最后变成一团乱麻。我这次把消息模型当作头等大事来设计,Topic、消息体、QoS逐一定清楚,后续所有开发都在这个约定上进行。
2.1 Topic设计:基于业务域的层级划分
MQTT的Topic是消息路由的核心,设计得好不好,直接影响后续扩展。我们按照“业务域/资源类型/资源ID/事件类型”的四层结构来划分:
code复制training/{classId}/student/{userId}/event
training/{classId}/teacher/{teacherId}/event
training/{classId}/ai/{aiSessionId}/stream
system/notice/all
system/monitor/{classId}/status
以课堂场景为例:
training/{classId}/student/{userId}/event:学员的个人事件通道,比如AI助教回复到达、个人答题批改结果。training/{classId}/teacher/{teacherId}/event:讲师通道,用于接收全班学员的实时数据聚合、系统预警。training/{classId}/ai/{aiSessionId}/stream:AI对话专用通道,流式输出走这里,避免大流量挤占普通事件通道。system/notice/all:平台级广播,全员在线通知、系统维护公告。
设计原则有三点:一是Topic数量和业务实体的对应关系要清晰,方便权限控制;二是通配符+层级,便于订阅端灵活监听;三是把AI流式数据单独划分Topic,防止高频数据阻塞关键事件。后面排查生产问题时,这个设计帮了大忙,直接按Topic过滤就能定位某条消息是否发出、落到哪个节点。
2.2 WebSocket帧结构与AI流式输出的“半包”问题
WebSocket侧我们重新定义了帧结构,不再让业务随意塞JSON字符串。统一的封装格式长这样:
json复制{
"type": "ai_stream",
"messageId": "a1b2c3d4-1234-5678-9abc-ef0123456789",
"timestamp": 1710000000000,
"payload": {
"sessionId": "ai_001",
"content": "这是一个流式片段",
"sequence": 15,
"finish": false
}
}
type字段标识消息类型,messageId用于去重和确认,sequence专门为AI流式输出设计。之所以单独加sequence,是因为AI生成的内容是分片推给前端的,AI模型推理速度和网络传输速度不一致时,客户端可能收到乱序的数据片段,没有序号就无法重组。
所谓“半包”问题在TCP层很常见,WebSocket虽然是消息帧协议,但消息体过大时,到了TCP层会被拆成多个包,网络中传输时可能乱序到达。大帧必须由服务端在内存中完整拼好再交给业务处理,如果设计时没有帧长度上限,一个恶意客户端发送超大帧,直接就能把服务端内存打垮。我们在网关层设置了单条消息最大64KB,超过即断开。
2.3 MQTT的QoS级别在生产环境的取舍
MQTT的QoS有三个等级,很多新手总是想“QoS越高越安全”,全选QoS 2,结果把Broker性能拖垮。我这边实际使用的原则是:
| 场景 | Topic示例 | QoS级别 | 理由 |
|---|---|---|---|
| AI流式对话内容 | training/{classId}/ai/{sessionId}/stream |
QoS 0 | 流式片段丢一帧可以靠序号发现并请求重发,不受影响 |
| 批量批改结果通知 | training/{classId}/student/{userId}/event |
QoS 1 | 结果不能丢,QoS 1保证到达,且性能开销可控 |
| 课堂上限操作指令 | training/{classId}/teacher/{teacherId}/event |
QoS 1 | 指令丢失会导致学员端状态不一致,至少送达一次 |
| 全员系统公告 | system/notice/all |
QoS 1 | 可容忍少量重复,但不能丢 |
QoS 2我直接没有采用。QoS 2需要四段确认握手,吞吐量下降明显,在生产环境中,QoS 1配合业务层的幂等处理已经足够。例如批改结果推送时,客户端拿到messageId后去重,重复消息直接丢弃,这不比QoS 2可靠吗?系统可靠性是“端到端”的,只靠协议层堆可靠性,成本高收益低。
注意:MQTT的QoS语义是“从发送端到Broker、从Broker到订阅端”的分段保证,不是端到端保证。客户端如果中途掉线重连,QoS 1最多保证消息不丢,并不保证不重复。业务层做幂等是必须的,别把QoS当免死金牌。
3. WebSocket服务端架构:从单连接到集群路由
WebSocket接入层是整个实时通讯的大门,也是最容易被高并发击穿的地方。我先从单连接的设计讲起,再讲集群后的路由方案。
3.1 连接管理:内存会话表与路由策略
单机时代,连接管理很简单,用一个ConcurrentHashMap把userId -> WebSocketSession存起来就行。到了集群阶段,问题立刻出现:学员A连在Node-1上,讲师B连在Node-2上,A发消息给B,Node-1根本不知道B在哪个节点。
我们第一版处理方案是“本地路由+Redis Pub/Sub转发”,结构如下:
- 每个WebSocket节点维护本地会话表
LocalSessionManager。 - 节点之间通过Redis Pub/Sub广播路由消息:某条消息的目标用户不在本地时,把消息丢到Redis指定频道,其他节点订阅频道后,查询本地会话表做真正下发。
- 会话注册、注销时,同步更新Redis里的全局用户节点映射表,保存
userId -> nodeId。
这套方案的好处是无需引入额外中间件,用已有的Redis就能支撑几千连接的集群路由。但问题也很明显:Redis Pub/Sub是“发后即忘”的,如果某个WebSocket节点宕机,转发中的消息会直接丢失,而且在超大规模集群下,每个节点都会收到全量广播,CPU浪费在无效消息上。后来我们将路由层替换成MQTT,正好补上了这个短板:WebSocket节点作为MQTT客户端订阅特定Topic,收到消息后推送本地连接,天然支持广播与点对点路由,还有QoS保证。这也是我在文中反复说“WebSocket和MQTT是互补”的一个重要原因。
3.2 心跳参数调优:为什么“30秒一次”不是最优解
心跳是WebSocket长连接的保命机制。早期我们直接用30秒一个心跳包,看起来没什么问题,压测时却发现服务端压力巨大——2万连接,每30秒每个连接一个心跳包,换算下来每秒约667个心跳请求,再加上业务消息,Node.js单进程CPU直接上60%。
调整思路是分级心跳:
- 空闲连接(10秒内没有业务消息)心跳间隔放宽到55秒。
- 活跃连接(正在AI对话或课堂互动)不需要额外心跳,业务消息本身就能证明连接存活。
- 服务端连续3次没收到心跳(即约165秒无消息),判定连接已死,主动关闭。
这个调整把心跳负载降低了近一半,而且降低了移动端耗电。判断“连接存活”不能死守固定间隔,要根据业务活跃度动态调整。
3.3 集群广播的三种方案对比
做在线课堂时,讲师端需要“全班广播”,比如“现在全员开始做题”。这个广播动作有三个实现路径:
| 方案 | 优点 | 缺点 |
|---|---|---|
| 直接遍历本节点所有连接 | 实现简单 | 只覆盖本节点,不是真广播 |
| Redis Pub/Sub广播 | 引入成本低,已有Redis即可 | 消息易丢,量大时Redis网络开销大 |
| MQTT Topic广播 | 天然一对多,支持QoS,Broker承担复制 | 多维护一套MQTT,架构复杂度上升 |
我最终选择了“WebSocket节点订阅MQTT Topic”的模式。讲师端发广播时,后端服务把消息publish到training/{classId}/teacher/command,所有在线WebSocket节点都在订阅这个Topic,收到后查询本地会话表,向连接在各自节点上的学员推送。这套方案将“端到端的实时性”和“服务端的可靠性”结合起来,广播语义由MQTT保证,连接推送由WebSocket负责。
3.4 Spring Boot/Netty集成中的实操细节
技术栈上,我们的AI培训系统后端主体是Spring Boot 3.x,WebSocket接入层使用了Netty而非Spring自带的WebSocket实现。原因很简单:Spring WebSocket基于Servlet容器,Tomcat的连接线程模型在高并发下表现一般;Netty基于NIO,线程模型可控,内存管理更精细,适合作为长连接接入网关。
集成Netty WebSocket时要注意几个点:
- Boss线程和Worker线程比例先按CPU核数设置,之后压测再调。
- WebSocketFrame的最大长度必须显式设置,防止OOM。
- 消息读写要区分业务线程池和IO线程,不能让业务逻辑阻塞IO线程。
- 鉴权放在
HttpRequestHandler中,在WebSocket握手前完成,不要等连接建立了再断开。
如果你用的是Spring生态,且不想引入Netty,也可以基于Spring WebSocket + STOMP做,但一旦连接数到2万以上,还是建议切换到Netty方案。
4. MQTT Broker选型与AI培训场景压测
MQTT Broker是整个消息分发的中枢,选型错了后面性能优化做得再好都白搭。我选型的时候对比了EMQX、Mosquitto、NATS JetStream,结合AI培训系统的场景做了几轮压测。
4.1 为什么选EMQX而不是Mosquitto
很多教程喜欢用Mosquitto,因为它轻量、配置简单,作为学习和小型项目没问题。但AI培训系统要支撑数万个WebSocket节点的MQTT订阅,还有内部微服务高频发布消息,Mosquitto的单进程性能撑不住大规模连接,管理界面也基本没有,运维排障全靠手动查日志。
我们最终选了EMQX,核心原因有四个:
- 10万级并发连接实测稳定,而且连接层是平滑扩展的。
- 内置消息仪表盘,Topic监控、消息速率、订阅关系一目了然,排障效率高。
- 支持共享订阅,多个消费者分摊Topic消息时配置简单。
- 规则引擎可以直接把Topic消息转发到HTTP、Kafka或数据库,省掉了我们中间再写一层数据桥接服务。
4.2 压测数据与容量评估标准
压测环境是3台8C16G云服务器,Broker单节点。用EMQX自带压测工具模拟客户端连接,同时并发发布AI流式消息和事件通知,关键数据如下:
| 指标 | 数值 | 说明 |
|---|---|---|
| 最大在线连接数 | 10万 | 达到后TCP连接不再增长,CPU约70% |
| 消息吞吐 | 约2.5万条/秒 | 每条消息256字节左右,QoS 1 |
| 平均延迟 | 8ms | P99约25ms,在网络条件正常的机房内 |
| 单主题订阅者数 | 2000 | 一个班级的讲师广播订阅,可稳定支撑 |
容量评估下来,3节点EMQX集群可以覆盖我们规划的5万WebSocket代理连接和1000个并发班级广播,同时内部服务产生的AI批改事件也有余量。压测过程中发现一个典型问题:大量短连接反复建立会频繁触发TCP TIME_WAIT,压测时客户端需要开启长连接复用。
4.3 开课峰值流量:应对“8点上课”的尖峰冲击
AI培训系统有一个和普通IM系统很不一样的地方:流量有明显的“课程时间表尖峰”。每天晚上8点课开始,大量学员同时上线,瞬间建立连接、订阅班级Topic,10秒内连接数能翻好几倍。
尖峰流量下最先出问题的是连接的鉴权和Topic订阅。几百个学员同时订阅同一个班级Topic,如果Broker的订阅路由表实现不够高效,会发生短暂的CPU毛刺。我们的应对措施:
- 客户端启动预热:App和Web端提前5分钟建立WebSocket和MQTT连接,在空闲期完成鉴权和Topic订阅,直播开始时不需要再做连接建立动作。
- 服务端对“批量订阅”做合并,客户端一次性传入多个Topic批量订阅,减少握手次数。
- 数据库侧的鉴权操作做Redis缓存,避免尖峰直接打到数据库。
这套“预热”策略非常有效,高峰期连接建立压力被削峰了60%以上。
4.4 离线消息与AI异步结果通知
MQTT的离线消息特性在AI培训系统里有两个典型的应用场景。
场景一是学员在弱网下断线,AI批改结果生成时学员并不在线。我们利用MQTT的持久会话(Clean Session=false),将离线消息暂存到Broker,学员重新上线后由客户端主动拉取未读消息。
场景二是异步AI任务。比如学员上传一段语音作业,AI语音评测服务需要处理很多秒,处理完成后服务端把结果publish到学员的事件Topic,学员端订阅Topic即可自动收到。这个链路完全异步化,通讯模块不需要知道AI处理服务在哪里、什么时候完成,只要双方约定同一个Topic和消息格式,就能解耦协作。
5. 性能优化清单:从千级连接到五万级连接的落地调整
性能优化不是某一个点的事情,我从系统层、框架层、协议层、客户端层四条线同时推进,这里给出一份可以直接抄作业的清单。
5.1 系统层调优:先看系统再怪代码
在线连接数到2万之后,首先扛不住的往往是Linux系统默认参数,而不是应用代码。我建议按下面的顺序检查:
文件描述符限制:
bash复制# 当前进程
ulimit -n 1048576
# 永久修改
echo "* soft nofile 1048576" >> /etc/security/limits.conf
echo "* hard nofile 1048576" >> /etc/security/limits.conf
TCP相关内核参数:
bash复制# 允许复用TIME_WAIT连接
net.ipv4.tcp_tw_reuse = 1
# 加快TIME_WAIT回收
net.ipv4.tcp_fin_timeout = 15
# 增大TCP接收/发送缓冲区
net.ipv4.tcp_rmem = 4096 87380 16777216
net.ipv4.tcp_wmem = 4096 65536 16777216
# 开启TCP KeepAlive
net.ipv4.tcp_keepalive_time = 600
net.ipv4.tcp_keepalive_intvl = 30
net.ipv4.tcp_keepalive_probes = 3
网络软中断绑定:如果机器网卡多队列,把网卡中断和应用程序CPU核心绑定,避免所有软中断挤在同一个CPU上。
这些参数调整完之后,再去排查应用代码,你会发现很多“并发问题”在系统层就已经解决了。
5.2 Netty与WebSocket服务端的线程模型优化
Netty默认的线程模型是BossGroup负责Accept连接,WorkerGroup负责IO读写。在线连接数大但消息量不均匀时,会出现部分Worker线程繁忙、部分空闲的情况。
我们优化时把业务处理和IO处理拆开:
- IO线程只做字节解码、帧拼接、心跳检测。
- 收到完整WebSocket消息后,交给独立的业务线程池处理,避免AI流式转发这类业务计算阻塞IO线程。
- 写回客户端时,通过
channel.writeAndFlush异步操作,不要同步等待写完成。
还有一个容易忽略的点:writeAndFlush在客户端消费不过来时,会积压在Channel的发送缓冲中。一定要监控ChannelOutboundBuffer的大小,超过阈值时主动断连或降级丢弃非关键消息,避免内存无限制增长导致OOM。
5.3 消息压缩与序列化:Protobuf还是JSON?
AI培训系统的消息体很大一部分是AI生成的文本片段,动不动一个批次几KB。如果走MQTT Broker转发,几万连接同时订阅,网络流量会非常大。我们做了消息压缩和序列化两轮优化。
第一轮是打开压缩。在WebSocket层采用permessage-deflate扩展(WebSocket压缩),MQTT侧对超过1KB的消息用zstd压缩再publish。实测文本类消息体积下降60%到75%,代价是CPU增加约10%,整体划算。
第二轮是序列化选型。业务消息原来全部用JSON,清晰但体积大、解析慢。我们做了一个折中方案:对外API保留JSON格式便于联调,内网节点之间的消息与AI流式数据改用Protobuf。只有内部流转的消息用Protobuf,避免所有业务方都要改序列化协议,改动成本低很多。
Protobuf对比JSON的实测数据:
| 消息类型 | JSON大小 | Protobuf大小 | 序列化耗时(JSON) | 序列化耗时(Protobuf) |
|---|---|---|---|---|
| AI流式片段 | 512B | 180B | 0.12ms | 0.03ms |
| 课堂事件 | 1.2KB | 380B | 0.31ms | 0.07ms |
5.4 客户端侧优化:Web端和移动端区别对待
不要只优化服务端,客户端才是连接数的源头。
Web端我们做的是:
- 页面不可见时(
document.hidden),暂停AI流式消息渲染,只做消息队列缓存,恢复可见时再刷新。这一步直接减少了大屏端超过一半的无效渲染。 - 浏览器单页应用里,WebSocket实例全局复用,避免路由切换时反复重连。
移动端更硬核一些:
- 网络类型切换(Wi-Fi切4G/5G)会强制断开TCP连接,iOS和Android都要监听网络状态,在恢复网络后主动重连WebSocket,而不是等心跳超时。
- 移动端屏幕旋转和App切后台也会触发WebSocket假死,需要配合生命周期做静默重连。
- 移动端建议使用MQTT over WebSocket 或原生MQTT,在弱网下的重连和消息补偿机制比裸WebSocket更省电,移动端的实时通知类消息走MQTT通道,只有会话型交互才走WebSocket通道。
在移动端上,MQTT的小字节头、可靠重连和离线消息机制是巨大的优势。这正是为什么即使有了WebSocket,也要引入MQTT——不同端的网络环境决定了协议选择的必要性。
6. 线上问题排查实录:断连、积压与广播风暴
再说几个典型的线上问题。这些问题不是从书上看来的,都是实际运维中踩过的坑,排查过程能帮大家建立直觉。
6.1 “stream disconnected before completion: failed to send websocket request: io”故障链分析
这个报错很典型,我们在Prometheus监控大屏上看到多个用户同时出现这个错误。字面上看是WebSocket请求在完成前流就被断开了,根本原因排查却经历了三步:
第一步,先怀疑是服务端主动断开。查了WebSocket服务端日志,发现没有任何主动断开记录,反而是客户端在握手后几秒内就断开连接。
第二步,查看TCP层。用tcpdump抓包,发现客户端发出WebSocket握手请求后,服务端有响应,但响应到达客户端时连接已经被客户端RST掉了。
第三步,定位到客户端侧。让客户端开发者查看浏览器日志,发现是页面在WebSocket握手完成前,发起了一个beforeunload事件导致页面卸载,WebSocket被浏览器强制关闭。原因是我们在前端做了路由懒加载,页面切换时WebSocket实例没被清理干净,旧连接还挂在全局状态里。
解决方案是:确保在页面路由切换前,显式调用ws.close()并清空引用;同时后端增加一条握手完成后的“客户端空消息确认”,如果握手后5秒内没有收到客户端任何消息,则主动断开。这样既减少了僵尸连接,也让类似问题更容易在日志中暴露。
6.2 断线重连风暴:一次Broker过载事故复盘
有一次我们发布WebSocket服务端新版本,批量重启节点,结果所有在线用户几乎同时掉线,又同时触发重连逻辑。短时间内几万个WebSocket连接同时涌向MQTT Broker订阅同一个班级Topic,Broker连接数瞬间飙升,CPU打满,消息延迟从8ms涨到2秒。
问题的根子在于重连策略缺少“随机退避”。所有客户端的重连间隔都是一致的3秒,重启后重连时间自然成了精确的同时风暴。
修复方案是:
- 客户端重连加入抖动退避算法,间隔从2秒到30秒随机分布。
- 服务端平滑重启,分批下线节点,每次只下10%的节点,等待重连完成后再操作下一批。
- 增加Broker侧的连接速率限制,超过阈值的新连接排队等待。
之后再做发布操作就没有出现过重连风暴。任何时候改动服务端节点,都要把“客户端重连风暴”的可能性考虑进去。
6.3 AI批改消息积压导致的结果延迟
有一次用户反馈AI批改结果迟迟不出来,查了一下发现MQTT消息队列里堆积了十几万条批改消息。奇怪的是,Broker的吞吐并不低,问题出在消费者速度上。
我们的批改结果消费者是一个Spring Boot服务,处理时需要调用AI批改接口做二次校验,每个消息处理耗时约500ms。正常情况下这个速度也能消化,但那天恰好另一个业务也在批量发送批改请求,消费者线程数没有动态扩容,导致消息积压。
解决思路是:
- 消费者服务的线程池大小从固定20改为基础大小20、最大100的弹性配置。
- 增加背压机制:当Broker积压消息超过阈值时,暂停非核心AI预热任务,优先保障批改结果下发。
- 将耗时操作(AI二次校验)从消费线程中拆出去,消费线程只负责把结果写入Redis,由另一个异步任务池去调用AI校验,消费速度大幅提升。
实时通讯的性能瓶颈很多时候不在通讯本身,而在“吃到消息之后做什么”。把消费逻辑里耗时的事情拆掉,积压问题就解决了。
6.4 监控与告警:我盯的六个核心指标
实时通讯模块的监控指标,每个人都会列一长串,但我实际运维中真正每天看的就是六个:
| 指标 | 健康阈值 | 异常处理 |
|---|---|---|
| WebSocket在线连接数 | 在预期范围内平稳 | 持续升高或断崖下跌都要查 |
| WebSocket连接建立成功率 | >99% | 小于99%大概率是系统参数或代码问题 |
| 消息端到端延迟(P99) | <100ms | 超过后查Broker积压和消费线程 |
| MQTT消息积压数 | 无明显持续增长 | 积压增长说明消费者速度不够 |
| 节点CPU/GC停顿 | CPU<75%,GC停顿<100ms | CPU高优先查序列化和压缩 |
| Topic订阅者数量 | 与班级在线人数匹配 | 订阅者数量对不上说明Topic订阅逻辑有问题 |
另外还要监控TCP层的状态量:ESTABLISHED、TIME_WAIT、CLOSE_WAIT。CLOSE_WAIT异常增多,基本上是有业务代码没关闭连接或没有正常释放WebSocketSession,这种连接泄漏靠应用日志很难发现,但TCP状态表一眼就能看出来。
我在处理“stream disconnected before completion”的时候,正是靠TCP连接状态和Broker消息积压两个指标同时提示异常,才快速定位到问题不在服务端,而在客户端的连接生命周期管理。
7. 一点经验总结:实时通讯不是“连上就行”
这次AI培训系统实时通讯模块重建,最深的体会是:实时通讯真正的门槛不是把连接建起来,而是把连接之后的“消息生命周期”管好。连接只是通道,消息才是业务。通道是WebSocket还是MQTT,只是工具差异;消息怎么定义、怎么路由、怎么保证可靠、怎么监控,才是架构师真正要花心思的地方。
如果你也要做类似系统,我建议先花一周时间梳理业务场景,把“谁在什么时候产生什么消息、谁需要收到、丢了会怎样”全部列清楚,再回头选协议和架构。不要一上来就套WebSocket,也不要因为MQTT“高级”就直接上,一切以业务模型为起点。
另外,实验环境再完整,也无法提前暴露所有问题。压测时一切正常,一上线被真实网络环境打回原形,这种事我见得太多了。建议做好监控和日志链路追踪,让每条消息都能找到它在系统中的生命周期轨迹。这样出了问题,你能快速回答“这条消息到哪了、是被谁吞了、还是根本就没发出来”,就已经赢了一半。
