在接手企业IM(即时通讯)这个需求之前,我一直以为SpringBoot里加个WebSocket就是调个API的事,真正动手做才发现,从“能连上”到“能稳定支持群聊、@提醒、消息回执”之间隔着一条巨大的鸿沟。这篇文章把我在开发一个包含群聊、@提醒、消息回执的轻量级企业IM系统时的完整思考、架构选型、代码实现和踩坑记录整理出来。项目技术栈锁定SpringBoot + WebSocket + STOMP协议,全程围绕真实业务需求展开,适合已经掌握SpringBoot基础、想往企业级通信方向深入的同学参考。
开头先说说这套组合到底解决了什么问题。原生WebSocket只解决“双向长连接”的问题,它本身不规定消息格式、不定义路由规则、也没有心跳和订阅机制。STOMP(Simple Text Oriented Messaging Protocol)是建立在WebSocket之上的消息协议,它把这些缺失的规范补齐了——用类似“发送消息帧、订阅目的地、服务端推送”这种简洁的文本帧格式,把一条条消息路由到对应的订阅者。SpringBoot把这两者整合得非常顺手,加上spring-boot-starter-websocket依赖,几乎是无缝衔接。再加上@SendTo、@SendToUser这些注解和SimpMessagingTemplate,足以支撑群聊、点对点私聊、强提醒这类功能。它不像Netty那样偏底层要自己搞定一切,但对大多数业务系统来说完全够用,而且相当稳。
1. 内容整体设计与思路拆解
1.1 为什么不是“裸用WebSocket”,而是套一层STOMP
当时团队里有人提出直接用原生WebSocket写,Spring的WebSocketHandler接口看起来也不复杂,配上TextWebSocketHandler就能收发文本消息。但稍微深入想了一下就发现坑不少:原生WebSocket没有“主题”概念,每一条消息都要自己做路由分发,前端也要手动维护连接、手动判断消息类型,连接管理、心跳、重连这些都要从零手写。本质上,原生WebSocket相当于一条裸管道,你往里倒什么、怎么倒、怎么分,全得自己定义,还容易各写各的,导致代码风格完全不一样。
STOMP就是在这个管道上加了一整套轻量级的“托运规则”。发送方要指定destination目标地址,服务端能根据这个地址做转发,客户端可以主动订阅某个地址实时接收推送,消息本身就带类型、头信息和正文结构。这种模式跟消息队列非常像,但又是即时通信。用生活化的比喻来说,WebSocket是电话线路,STOMP是电话里的语言规范——大家都说同一种语言,听的人才能听懂对方在讲什么。更关键的是,Spring对STOMP有深度的原生支持,几乎不需要再封装底层通信逻辑,注意力能全部集中在业务上。
既然定了STOMP,服务端实现层的第一选择就是spring-boot-starter-websocket。这个starter依赖其实是把Spring MVC里的WebSocket支持和STOMP协议支持打包好了,类名都叫WebSocketMessageBrokerConfigurer,配置它来注册STOMP端点、配置消息代理和订阅前缀。这个配置类核心就两个方法:registerStompEndpoints用于注册前端连接的地址,configureMessageBroker用于设定消息的前缀路由规则。
1.2 三个核心需求背后的技术拆解
先说群聊。群聊表面上看起来就是“A发消息,群里所有人都能收到”,实现时涉及到一个关键问题——消息要不要先进数据库,再推给在线用户。如果只追求实时性直接推,用户不在线就收不到;如果先存库再推,又面临推送失败怎么办的补偿问题。最终我的设计是“存储和推送双轨并行”:把消息先持久化到MySQL,保证数据不丢,再通过STOMP的广播地址实时推给在线用户,不在线的人也能在下次登录时拉到历史消息。存储、推送各司其职,不互相拖累。
然后是@提醒。这块比看上去复杂,它不只是文本里带一个“@小明”,你需要在发送时就识别出被@的人是谁,还要在服务端把提醒处理成独立的业务记录。做的时候要区分两类提醒:站内信式的离线提醒(小红点、通知列表、未读消息数)和WebSocket的实时推送提醒。前者靠数据库表记录,后者靠STOMP的点对点推送。两条链路缺一不可。如果只做实时推送,那离线用户永远看不到有人@过他;如果只存数据库,那在线用户又不会实时弹出提醒。这两个手段必须同时做,并且要保证它们走的是同一套数据源。
消息回执是这三个里最有挑战性的。它的本质是“某个用户看到了某条消息”的状态同步。要实现它,第一件事是明确数据模型:消息表、会话表、用户消息状态表各自承担什么职责,量级怎么预估。第二件事是解决已读状态的推送时机:是每次读一条推一条,还是批量累计后统一推?前端如果每读一条消息就推送一条回执,服务端压力很大;如果批量推,前端UI的实时性又会下降。我的做法是在前端做一个已读回执的“合并上报”——用户停留在某个会话里3秒以上,把这个会话里所有未读消息标记为已读,一次性上报给服务端,服务端再广播给该会话的其他成员,把一个一个的小回执合并成一条,复杂度瞬间就降下来了。
1.3 为什么选择STOMP的心跳机制而不自己写
WebSocket本身确实有Ping/Pong帧,但前端浏览器里的WebSocket API偏偏不开放直接发送Ping帧的能力,只能被动接收。这意味着真要实现长连接的保持,还是得在应用层自己做心跳,对多数团队来说就是自己写定时器。STOMP协议干脆把心跳定义到了协议层:客户端和服务端在CONNECT帧里通过heart-beat头协商心跳间隔,之后双方按照协商的结果定时发送一个空行(EOL)或合法帧来保持连接活跃。Spring的WebSocketStompClient自带心跳支持,服务端也内置了心跳协商机制,配置完成后几乎不用写多余的代码。代理层面心跳断了,Spring会自动触发连接关闭和清理逻辑,大大减少了写“僵尸连接”清理代码的工作量。
2. 核心细节解析与实操要点
2.1 项目依赖与基础配置
从依赖入手,pom.xml里核心要加的是spring-boot-starter-websocket。如果你在做安全认证,还要考虑spring-security整合的问题,但初版建议先把安全和WebSocket解耦,单独拎出来做。项目用的是SpringBoot 2.7.x版本,JDK 8或11都行。还有一个容易踩的坑是SpringBoot版本和Spring框架版本之间的API差异,比如3.x之后很多类改了包名或废弃了旧方法,网上很多教程是2.x的写法,照搬到3.x会直接编译报错。
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-websocket</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-redis</artifactId> </dependency>Redis在这里不是装饰品,它用于在线用户状态管理和未读消息计数。如果部署在集群环境,Redis还是广播消息的关键一环,这个后面会专门说。数据库用MySQL即可,表结构设计是重点。
配置文件里,核心的不是端口,而是STOMP的消息前缀和端点路径规划。我将应用内消息代理(Simple Messaging Broker)的“目的地前缀”设为/app,表示客户端发送到服务端的消息都会经过这里;“订阅前缀”设为/topic和/queue,分别表示广播订阅和点对点订阅;服务端点路径设为/ws-im,表示前端连接WebSocket时的入口。这三个路径设计看似简单,实际上决定了后续代码风格的统一性,最好一开始就规划好。
server: port: 8080 spring: application: name: im-server redis: host: 127.0.0.1 port: 63792.2 WebSocketConfig配置类与STOMP端点注册
WebSocket配置类要继承WebSocketMessageBrokerConfigurer并实现configureMessageBroker和registerStompEndpoints两个方法。注册端点这里有个细节:setAllowedOriginPatterns。SpringBoot 2.7之后的版本对跨域限制很严,如果前端和你的服务端不在同一域名下,用setAllowedOrigins("")经常会被浏览器拦下,而setAllowedOriginPatterns("")则能正确匹配任意来源。因为生产环境前端大概率部署在独立域名,所以必须用后者。
@Configuration @EnableWebSocketMessageBroker public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { @Override public void registerStompEndpoints(StompEndpointRegistry registry) { registry.addEndpoint("/ws-im") .setAllowedOriginPatterns("*") .withSockJS(); } @Override public void configureMessageBroker(MessageBrokerRegistry registry) { registry.enableSimpleBroker("/topic", "/queue"); registry.setApplicationDestinationPrefixes("/app"); registry.setUserDestinationPrefix("/user"); } }这个代码里最容易被忽略的是.sockJS()。SockJS是一个浏览器端降级方案,它让不支持WebSocket的老旧浏览器自动回退到HTTP长轮询或Server-Sent Events。虽然现在主流浏览器基本都支持WebSocket,但企业内部用户偶尔会用到比较老的内网浏览器,保留SockJS基本是零成本的保险丝。前端对应要引入sockjs-client库,后面我会给出完整的前端代码。
configureMessageBroker里的setUserDestinationPrefix("/user")也值得细说。它在点对点通信中扮演独特角色。当你给某个具体用户推送消息时可以这样写:convertAndSendToUser(userId, "/queue/reply", payload),此时真正推送的目的地会被自动转换成/user/{userId}/queue/reply。这种转换规则就是serveUserDestinationPrefix决定的,前端订阅时也用/user/queue/reply这个地址来接收。基于这套机制,@提醒的异步推送、私聊消息的实时送达就都能落到确切的个人头上。
2.3 用户身份与WebSocket Session绑定
WebSocket连接建立时,前端会传一个token,服务端需要解析token并把它和WebSocket的Session关联起来。很多新手喜欢在消息处理方法里解析参数判断是谁发的,但正确的做法是利用ChannelInterceptor在握手阶段就完成身份绑定。在Spring的STOMP协议栈中,CONNECT帧到达时可以拦截Connection。在preSend里获取StompHeaderAccessor,取出token并解析用户,然后把用户信息放进atttributes。后续消息处理时,再从headerAccessor的sessionAttributes里取出用户信息,整个过程是线程安全的。
@Component public class UserChannelInterceptor implements ChannelInterceptor { @Override public Message<?> preSend(Message<?> message, MessageChannel channel) { StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class); if (accessor != null && StompCommand.CONNECT.equals(accessor.getCommand())) { String token = accessor.getFirstNativeHeader("Authorization"); if (token != null) { Integer userId = parseUserIdFromToken(token); if (userId != null) { accessor.setUser(new StompPrincipal(String.valueOf(userId))); } } } return message; } }StompPrincipal实现了java.security.Principal接口,只有getName方法。这个Principal天然就是用户身份的来源,@SendToUser和convertAndSendToUser底层也都是通过它来识别“发给哪个用户”。如果不做绑定而改用param传参,那一条消息就能冒充其他人发出去,安全直接就崩了。这里我把解析token的逻辑只写了个注释位置,生产环境建议接入Spring Security的认证过滤器或直接用JWT工具类解析。
2.4 消息结构体的设计
消息体设计直接影响前后端联调效率。直白说,消息不仅要携带内容文本,还需要带上发送者信息、消息类型、接收范围、时间戳、客户端生成的消息ID等元数据。我常用的DTO结构是这个样子的。
{ "type": "CHAT", "conversationId": 1001, "senderId": 1, "senderName": "张工", "content": "大家好,这个需求我看过了", "timestamp": 1689058800000, "clientMsgId": "uuid-xxxxx", "mentionedUserIds": [2, 3] }clientMsgId是前端生成的唯一标识,这个字段极其有用。IM系统最容易出现的就是消息重复——断网重连导致消息重发、服务端重试导致重复入库。前端拿到clientMsgId可以做幂等去重,服务端也能用它做唯一索引来辅助防重。我的做法是在消息表中对client_msg_id加唯一索引,一旦检测到重复就拒绝写入并返回“重复消息”状态码,保证一个客户端消息只落库一次。
3. 实操过程与核心环节实现
3.1 数据库表结构设计与理由
实现三个核心功能,至少需要四张表:用户表、会话表、会话成员表、消息表。如果要做@提醒和已读回执,消息表和会话成员表还需要额外加字段或关联表。完整建表语句如下。
CREATE TABLE im_user ( id BIGINT PRIMARY KEY AUTO_INCREMENT, username VARCHAR(64) UNIQUE NOT NULL, display_name VARCHAR(64) NOT NULL, avatar_url VARCHAR(255), status TINYINT DEFAULT 1, created_at DATETIME DEFAULT CURRENT_TIMESTAMP ); CREATE TABLE im_conversation ( id BIGINT PRIMARY KEY AUTO_INCREMENT, name VARCHAR(128), type TINYINT COMMENT '1-单聊 2-群聊', owner_id BIGINT, last_message VARCHAR(512), last_message_at DATETIME, created_at DATETIME DEFAULT CURRENT_TIMESTAMP ); CREATE TABLE im_conversation_member ( id BIGINT PRIMARY KEY AUTO_INCREMENT, conversation_id BIGINT NOT NULL, user_id BIGINT NOT NULL, unread_count INT DEFAULT 0, last_read_message_id BIGINT DEFAULT 0, muted TINYINT DEFAULT 0, UNIQUE KEY uk_conversation_user (conversation_id, user_id) ); CREATE TABLE im_message ( id BIGINT PRIMARY KEY AUTO_INCREMENT, conversation_id BIGINT NOT NULL, sender_id BIGINT NOT NULL, msg_type TINYINT COMMENT '1-文本 2-图片 3-文件 4-系统', content TEXT, client_msg_id VARCHAR(64) UNIQUE, mentioned_user_ids VARCHAR(255) COMMENT '冗余字段,逗号分隔', created_at DATETIME DEFAULT CURRENT_TIMESTAMP, KEY idx_conversation_id_time (conversation_id, created_at) );im_conversation_member表是整个系统的关键,它同时承担了三个职责:一是记录某个用户在某个会话里的未读数;二是记录某个用户在该会话中的已读位置;三是存储会话提醒设置。unread_count的更新发生在两个时机:用户不在线或者在线但没打开会话时,收到新消息则增加计数;会话被打开且用户停留时,清理对应会话的未读。last_read_message_id用于计算回执的拉取差量,它能告诉我哪些消息是这个用户没读过的。
消息表里特意冗余了mentioned_user_ids字段,这个字段存储本消息@了哪些用户,用逗号分隔而不是搞一张关联表。原因很简单:提醒我们要关心的是“这条提醒是否已读”,是否@了谁本来就是消息的一部分,拆成关联表反而增加查询负担,而且IM消息的@通常不需要事后做复杂的维表分析。如果在做更多关联数据统计分析,再考虑拆表。
3.2 群聊消息的完整处理链路
服务端接收群聊消息的入口是一个Controller类。前端把消息POST到/app/chat/send这个STOMP目的地,实际上就是发到服务端,服务端做四步操作:保存消息到数据库、更新会话列表里的最后一条消息摘要、通过convertAndSend广播给/topic/chat/{conversationId}、处理@提醒和未读计数。其中广播给群聊的两个动作有先后——先存库再广播,顺序绝对不能反。如果先广播后存库,一旦存库失败(数据库临时故障),在线用户看到了消息、离线用户又拉不到,数据就永久不一致了。
@MessageMapping("/chat/send") public void handleGroupChat(@Payload ChatMessage chatMessage, SimpMessageHeaderAccessor accessor) { Integer senderId = ((StompPrincipal) accessor.getUser()).getNameAsInt(); chatMessage.setSenderId(senderId); chatMessage.setTimestamp(System.currentTimeMillis()); Long messageId = messageService.saveMessage(chatMessage); chatMessage.setId(messageId); conversationService.updateLastMessage(chatMessage.getConversationId(), chatMessage.getContent()); if (chatMessage.getMentionedUserIds() != null && !chatMessage.getMentionedUserIds().isEmpty()) { mentionService.processMention(chatMessage); } messagingTemplate.convertAndSend("/topic/chat/" + chatMessage.getConversationId(), chatMessage); }@MessageMapping是Spring处理STOMP客户端消息的核心注解。客户端发送到/app/chat/send的消息,最终会映射到handleGroupChat方法。这里有个程序员容易迷惑的地方:方法返回void,因为实际是通过messagingTemplate手动推送的,而不是利用@SendTo注解自动转发。两者都能实现广播,但手动推送更灵活,因为可以在推送前做任意前置业务操作,比如这里就做了存库和@提醒处理。
群里核心逻辑里要留意消息广播时谁应该收到。如果群成员里有用户已经把所有消息都设置成免打扰,广播前需要查一下该成员的免打扰状态;但也不要每次广播都全群扫描一遍成员表。我的方案是群成员状态改变时做缓存标记,广播时只筛选出需要推送的成员列表。小程序和App端推送触达是企业IM必不可少的需求,但这里先只在实时推送层面处理,推送服务单独拆出去做。
3.3 私聊与点对点消息的STOMP实现
虽然标题主打群聊,但@提醒本质上是点对点推送,私聊场景在实际企业IM里也绕不开。私聊的双人会话就是一个小群,带宽上的瓶颈远小于群聊,但实现上有自己的独特问题:私聊时消息应该只发送给会话里的两个人,不能广播到全局。使用convertAndSendToUser可以严格做到按用户隔离推送。
@MessageMapping("/chat.private") public void handlePrivateChat(@Payload ChatMessage chatMessage, Principal principal) { Long conversationId = chatMessage.getConversationId(); Integer senderId = Integer.parseInt(principal.getName()); chatMessage.setSenderId(senderId); messageService.saveMessage(chatMessage); conversationService.updateLastMessage(conversationId, chatMessage.getContent()); Integer receiverId = chatMessage.getReceiverId(); messagingTemplate.convertAndSendToUser(String.valueOf(receiverId), "/queue/private", chatMessage); messagingTemplate.convertAndSendToUser(String.valueOf(senderId), "/queue/private", chatMessage); }注意convertAndSendToUser的目标地址前缀,前端订阅时需要订阅/user/queue/private。此时用户自身的身份标识已经由Principal对象提供,不需要前端再传senderId,防止伪造。为了保险起见,服务端收到消息后还是要校验当前登录用户是否真的是会话成员,不是成员直接抛出异常拒绝投递。这个坑我见过不止一次,有的项目在私聊消息里带上receiverId,但没有校验receiverId和会话成员的关系,结果通过构造请求就能给任何人发私信。
3.4 @提醒的异步处理和离线补偿
@提醒的处理我拆成了两个模块:实时推送和离线存储。实时推送是在广播群聊消息的同时,对被@的用户单独发一条STOMP点对点消息,前端收到提醒消息后展示系统通知或者弹气泡。离线存储是在数据库里为每个用户维护一张提醒表,记录“谁在某群聊里提到了我”,当用户上线时拉取未读提醒。简化的提醒表结构可以复用im_message里的mentioned_user_ids字段,但真要做好还是需要独立的提醒查询接口。
离线补偿的核心在于会话列表页和@提醒的高度重复。当用户上线后,前端通常第一件事就是拉取会话列表,然后发现某个会话未读数很高,点进去才知道自己被@了。这个体验不够明显,尤其当未读数很多时。理想的产品设计是@提醒单独列一栏,有未读的@要红色高亮提示。所以独立提醒表设计成下面这样。
CREATE TABLE im_mention ( id BIGINT PRIMARY KEY AUTO_INCREMENT, message_id BIGINT NOT NULL, conversation_id BIGINT NOT NULL, mentioned_user_id BIGINT NOT NULL, read_flag TINYINT DEFAULT 0, created_at DATETIME DEFAULT CURRENT_TIMESTAMP, KEY idx_user_read (mentioned_user_id, read_flag) );每次有人发消息时如果带了mentionedUserIds,服务端开启异步线程批量插入提醒记录,同时通过convertAndSendToUser发送实时推送提醒。异步处理直接用Spring的@Async注解配合线程池,核心逻辑不要阻塞消息广播的主链路。生产环境如果消息量很大,可以用MQ异步解耦,但阿里的RocketMQ、RabbitMQ的引入会增加系统复杂度,第一版完全可以用线程池做异步落库,等消息量上来再考虑MQ。
3.5 消息回执的实现思路与数据流转
消息回执这种需求,数值上分两种:一种是群聊里进来新消息后,其他成员是否看到了;另一种是私聊里对方是否已读。本质上都是同一套机制,用im_conversation_member表的last_read_message_id和im_message表里的消息id做差量比对。只要last_read_message_id >= 某条消息的id,就说明这条消息已读。这是IM系统里常见的“水位线”方案,跟数据库binlog同步的位点概念很类似。
前端什么时候上报已读?我做的方案是在会话页面里添加观察者机制:当某个会话进入激活状态且停留时间超过3秒后,前端拿到该会话当前最新的消息ID,调用STOMP服务端接口上报“已读到这条消息”。服务端更新im_conversation_member表的last_read_message_id,并把unread_count清零,然后向这个会话的其他成员广播一条回执通知。
@MessageMapping("/chat.read") public void handleRead(@Payload ReadReceipt receipt, Principal principal) { Integer userId = Integer.parseInt(principal.getName()); conversationService.markRead(userId, receipt.getConversationId(), receipt.getLastReadMessageId()); ReadReceiptNotify notify = new ReadReceiptNotify(); notify.setConversationId(receipt.getConversationId()); notify.setUserId(userId); notify.setLastReadMessageId(receipt.getLastReadMessageId()); messagingTemplate.convertAndSend("/topic/read/" + receipt.getConversationId(), notify); }消息回执的难点在于已读状态更新后,其他成员界面上的“已读1人”、“已读2人”要实时变化。比如群里五个人,A发了消息,B读了之后C的界面就要显示“B已读”。这依赖上述广播机制,但广播给所有成员显然没考虑那些不在会话里的人。后来我做了优化:广播回执前检查在线用户列表,只给在线的会话成员推送,不在线的等下次上线时拉取会话状态时自动同步。这种优化能省很多无效推送,尤其当群成员很多时。
我踩过的一个大坑是回执的乱序问题。前端上报已读是有可能并行到达的,服务端如果只按收到的顺序更新水位线,后到的低ID可能覆盖先到的高ID,导致已读位置倒退。解决办法很简单:更新时加一个判断,只有新上报的ID大于当前last_read_message_id时才更新。
UPDATE im_conversation_member SET last_read_message_id = ?, unread_count = 0 WHERE conversation_id = ? AND user_id = ? AND last_read_message_id < ?这句SQL的where条件就是矛与盾的结合,如果新水位线比旧的小,UPDATE影响行数为0,不会污染数据。
3.6 前端关键代码与联调技巧
即使后端逻辑完美,前后端联调才是真正耗费心力的环节。我先给了一套基于stompjs和sockjs-client的前端基础代码,协议栈顺序是先SockJS再重新包装成STOMP客户端。这里直接以现代浏览器的原生WebSocket连接方式为例,因为用SockJS时浏览器地址会变得很绕,不好直接分析。
import SockJS from 'sockjs-client'; import { Client } from '@stomp/stompjs'; const client = new Client({ webSocketFactory: () => new SockJS('http://localhost:8080/ws-im'), connectHeaders: { Authorization: 'Bearer ' + localStorage.getItem('token') }, reconnectDelay: 5000, heartbeatIncoming: 4000, heartbeatOutgoing: 4000 }); client.onConnect = () => { // 订阅群聊频道 client.subscribe('/topic/chat/1001', (message) => { console.log('收到群聊消息', JSON.parse(message.body)); }); // 订阅个人提醒 client.subscribe('/user/queue/private', (message) => { console.log('收到私聊或提醒', JSON.parse(message.body)); }); }; client.activate();connectHeaders就是之前服务端UserChannelInterceptor里取token的Header,设置heartbeatIncoming和heartbeatOutgoing各4000毫秒,意思是客户端每4秒发一个心跳、服务端每4秒也发一个心跳。这个间隔不是随意定的,太短会造成无效的网络包频繁冲击,太长则不能及时发现死连接。实际运营经验来看3到5秒比较合适,太长会让Nginx等代理层的空闲超时把连接回收掉。
联调时初始阶段最容易遇到404问题,症状是WebSocket连着连着突然跨域错误或者握手失败。遇到404先检查服务端端点路径是否一致——我见过前端连/ws、后端注册/ws-im、中间没有映射的情况,也见过端口不一致导致所有请求打到别的服务上的情况。接着检查SpringBoot的日志,如果看到“Failed to handshake”多半是allowedOriginPatterns没写对,或者是Spring Security拦截掉了CONNECT帧。建议联调时后端先把Spring Security关了,等所有消息通道走通后再开安全拦截器。
4. 常见问题与排查技巧实录
4.1 连上就断,心跳之间的拉锯战
接入之后第一个高频反馈是“刚连上就断了”,或者“每隔几分钟就掉线重连”。这类问题排查时套路是有顺序的,先看客户端日志有没有Reconnect字样,再看服务端日志有没有Connection closed,最后再看代理层日志有没有Idle timeout。前端环境里这个体验最明显:WebSocket连接由浏览器发起,Nginx接收后转发给后端。Nginx默认的proxy_read_timeout是60秒,如果后端在一分钟内没有任何数据返回,Nginx就会主动断开连接。而群聊不是每分钟都有消息的,如果没有心跳机制,60秒一到就被掐断。
我的解决方式是在Nginx配置里同时增加proxy_read_timeout和proxy_send_timeout到600秒,并开启WebSocket升级相关的Header。
location /ws-im/ { proxy_pass http://localhost:8080; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; proxy_read_timeout 600s; proxy_send_timeout 600s; }即使没有Nginx,直连SpringBoot时也会因为装载了STOMP心跳而自动协商保持连接,但生产环境基本不可能不用Nginx,所以这一段是必备配置。还有一点要注意的是,SpringBoot应用本身对WebSocket设置session超时时间。stomp session默认值是60秒还是10分钟会根据版本变化,建议在WebSocketConfig里显式配置。
registry.addEndpoint("/ws-im") .setAllowedOriginPatterns("*") .setHandshakeHandler(new DefaultHandshakeHandler()) .withSockJS() .setDisconnectDelay(30 * 1000);setDisconnectDelay只对SockJS生效,控制SockJS连接僵死时的关闭延迟。了解就行,它通常不是断线主因。真正的主因百分之八十是代理超时,优先排查Nginx。
4.2 广播风暴:一条消息被群成员重复收到
有次压测时发现,500人的大群发一条消息,数据库和消息中间件压力骤增,频繁出现重复日志。排查出来是前端订阅姿势不对。前端把某个群聊的通知和提醒都订阅了一遍,还可能在多个组件里重复创建STOMP连接。统计发现一条群聊消息被同一个用户接收了3次甚至5次,群成员数一乘,服务端压力就爆炸了。
处理方法是前端规范化STOMP连接管理,整个应用只维护一个全局STOMP客户端,所有页面组件通过单例或状态管理共享,不要每次进入会话页面都new一个Client。订阅方面只保留一个订阅回调地址,其他组件用事件总线分发消息。服务端防止重复推送的逻辑也值得做:广播前检查消息ID,同一处理的同一个topic订阅最多产生一条消息。但真正根因还是前端,这里提出来给做前端联调的同学排雷。
4.3 在线用户状态显示不准确
界面上看到的在线状态不靠谱,一会儿在线一会儿离线,白名单同事的体验是“他明明在电脑前面,头像却是灰的”。即时通讯系统的在线状态本质上就是依赖心跳和会话的活跃度。实际排查发现,很多人把浏览器切到后台Tab页时,浏览器会暂停定时器执行,心跳的发送也被暂停了。如果在后台停留时间超过服务端心跳检测阈值,服务端会判定连接超时并关闭会话,用户再切回来界面就显示离线了。
这个很难从服务端根上解决,用STOMP心跳本身就依赖浏览器JS的运行,Tab在后台被冻结时JS根本没机会跑。客户端侧的替代方案是:监听页面的visibilitychange事件,页面变为可见时立刻手动调用client.deactivate()再activate()重新建立连接,跟手机App退后台断网再回前台重连的逻辑一样。在线状态的最终判定建议设置宽松一些,服务端10秒没有心跳才标记离线,而前端60秒内有过交互都算在线。具体业务具体调整,但别把超时设得太短,企微钉钉这类系统的在线状态也不可能是实时秒级切换的,很多都有几十秒的延迟容忍。
4.4 水平扩展时消息广播如何不重复不丢失
单机部署一切顺利,系统要横向扩容压测时,问题就来了:两台应用实例都在运行,用WebSocket连接的用户可能连着实例A也可能连着实例B。如果某个用户在群聊里发了一条消息,这条消息只存到了数据库,广播时却只调用了本机的convertAndSend,那么连着实例B的用户完全收不到实时推送。这就是典型的“广播只在单机可见”的问题。
要解决它,依赖一个全局的消息通道。Spring的STOMP消息代理功能里,启用外部Broker(比如RabbitMQ)能解决这个问题,但配置复杂度直接上一个台阶。轻量级方案则是用RedisPubSub做消息转发:每台应用实例都订阅同一个Redis频道,当A实例收到群聊消息并落库后,把消息发布到Redis,同时所有实例都订阅这个Redis频道,B实例收到后再把它转换成STOMP推送发给自己节点上在线的订阅者。实现上通过RedisMessageListenerContainer监听频道,在onMessage里调用SimpMessagingTemplate.convertAndSend推给本机连接。
@Component public class RedisMessageSubscriber implements MessageListener { @Autowired private SimpMessagingTemplate messagingTemplate; @Override public void onMessage(Message message, byte[] pattern) { String payload = new String(message.getBody(), StandardCharsets.UTF_8); RedisMessageDTO dto = JSON.parseObject(payload, RedisMessageDTO.class); if ("CHAT".equals(dto.getType())) { messagingTemplate.convertAndSend("/topic/chat/" + dto.getConversationId(), dto.getData()); } } }单机模式下不需要Redis发布订阅,因为本机的convertAndSend自己就能广播给本机连接,但为了以后扩容,我建议从一开始就统一走Redis分发,这样扩容时代码零改动。代价仅仅是每条消息多一次Redis发布订阅的开销,对企业IM系统来说完全可以接受。
4.5 会话列表的SQL性能与未读数即时更新
会话列表页是IM应用里访问频率最高的接口之一,几乎每次App启动、每次进入首页都要拉取一次。初次实现时直接查im_conversation_member表再join im_conversation,再做子查询算最后一条消息,用户量一旦过万就明显变慢。MySQL分析后慢查询SQL集中在orderby和groupby,会话列表接口平均耗时300多毫秒,不可接受。
优化策略是先把im_conversation表里的last_message和last_message_at冗余字段用起来。会话列表只需要查询当前用户所属的会话ID、名称、最后一条消息内容和时间,不再需要实时去消息表里count和max。会话状态(未读数、置顶、免打扰)在im_conversation_member表里查,再加上start index。SQL拆分后复杂度大幅下降。
SELECT c.id, c.name, c.type, m.unread_count, m.last_read_message_id, c.last_message, c.last_message_at FROM im_conversation_member m INNER JOIN im_conversation c ON m.conversation_id = c.id WHERE m.user_id = ? ORDER BY c.last_message_at DESC LIMIT 20;这看起来平淡无奇,但结合未读数的快速更新能解决用户最常吐槽的“我看过消息了,角标还在”问题。前文消息回执里我们说已读上报时会更新水位线,会话列表要实现即时刷新,就必须在更新回执后把未读计数清零,再通过Redis或直接通过会话列表接口套上时间戳下发。如果用户停留在会话列表页,前端需要由STOMP主动推动一条“会话状态更新事件”来驱动刷新,这也是我为什么把回执广播到/topic/read/{conversationId}的原因——这个topic可以触发前端各个页面的状态联动刷新。
5. 离在线消息与消息补偿机制的补充
5.1 离线消息拉取策略
用户重新上线后,除了订阅实时推送,还需要拉取离线期间的消息。这是IM系统的基本要求。实现方式简单直接:客户端登录成功后向服务端请求“增量消息拉取”接口。增量消息的判断基础就是im_conversation_member表里的last_read_message_id。用户在线期间实时推送的消息实时更新,离线期间的消息则用这个ID去消息表拉取。
public List<ChatMessage> pullOfflineMessages(Integer userId, Long conversationId) { ConversationMember member = conversationMemberMapper.selectByUserAndConversation(userId, conversationId); if (member == null) { return new ArrayList<>(); } return messageMapper.selectMessagesAfter(conversationId, member.getLastReadMessageId()); }批量拉取时按会话分组拉取,不要一条一条请求。假如用户离线期间有10个会话各产生了20条消息,一次拉取全部会话的全部增量消息,HTTP接口返回一个Map<conversationId, List >结构,前端再分发给各路组件。这个设计从一开始就要搭好,否则后面增加群人数和消息频率时会反复改接口结构。
拉取完增量消息后,前端要把最后一条消息的ID作为nextCursor存下来,下次再拉取就直接增量。这跟翻页的思路类似,但是状态是连续增长的。此时必须在服务端做幂等,防止前端重复拉取同一条数据导致重复渲染。
5.2 重连期间消息补偿与幂等
WebSocket连接不是永远稳定的,用户在地铁上、网络切换过程中都会触发断线重连。重连期间如果恰好有群聊消息被实时广播,那这个用户就漏掉了。因此重连成功后,前端需要做一次主动拉取消息的动作,保证漏掉的消息被补偿回来。最简单有效的做法是:每次STOMP客户端onConnect成功之后,调用一次增量拉取接口。这个动作是幂等的,服务端会根据last_read_message_id过滤,不会重复返回已读消息。
复杂点在于消息实时推送和增量拉取之间可能产生交叉重复——推送刚发出来还没到达前端,增量拉取已经返回了包含这条消息的数据。所以前端要把每条消息里的clientMsgId维护成一个Set,重复消息直接丢弃。这个去重机制前面消息结构体设计时提到过,这里是它最重要的应用场景。没有这套去重,断线重连一次就重复展示一遍,用户体验极其糟糕。
5.3 消息已读回执的批量SQL优化
群里聊天,每个人都读了消息后要给发送者或群成员推送“XX已读”状态,如果每条已读都单独UPDATE一次,高并发下数据库扛不住。实测一个200人群,如果200人几乎同时上线并读取新消息,数据库会瞬间产生几百上千条UPDATE。我把已读上报的SQL改成批量模式,前端上报时不只上报一个lastReadMessageId,而是上报一个“已读到某时间点之前的所有消息”时间戳,服务端按会话把所有早于该时间戳的未读消息一次性标记已读。这个方法屡试不爽。
更进一步,用Redis的Hash结构缓存会话的已读水位线。例如key为conversation:read:{conversationId},field为userId,value为已读位置。更新时先写Redis,再由定时任务异步刷新到MySQL。这样在线用户的已读状态变化非常快,而数据库的压力能平摊到秒级批量提交。不过这个方案引进了Redis和MySQL的一致性维护问题,需要接受Redis丢了可以从MySQL兜底重建。企业IM场景可以接受这个策略。
6. 部署、安全以及上线前的注意事项
6.1 从开发到生产:Nginx、端口和内网穿透
开发时localhost直连很省事,但生产环境部署涉及到的网络拓扑比想象中复杂。最常见的部署形态是Nginx位于前端,统一暴露443或者自定义端口,WebSocket请求通过Nginx反向代理到内网的后端端口。Nginx必须开启Upgrade头,还要处理WebSocket的长期存活问题。另一个坑是容器化部署时Docker的端口映射。docker run时只映射了8080端口,但WebSocket连接是从Nginx容器转发到后端容器的,端口如果没配对好也会出现连不上。我建议部署时先不用Docker网络,直接用宿主机端口验证连通性,通了再考虑容器编排。
SpringBoot应用自己也可以嵌入一个WebSocket端点,端口和HTTP端口相同。如果同时开启了Tomcat的HTTPS和HTTP,WebSocket会强制走HTTPS升级,此时浏览器连接也要用wss协议前缀。生产的WebSocket地址形态一般是wss://im.example.com/ws-im,wss对应TLS加密,和https一样需要证书。证书过期、域名不匹配都会导致握手失败,这个问题非常容易出现在联调阶段,建议上线前先检查证书链是否完整。
6.2 认证与权限控制
STOMP协议本身不提供认证能力,认证的职责完全落在服务端。有两种做法:第一种是连接时通过CONNECT帧的Header携带token,用ChannelInterceptor解析;第二种是先走HTTP登录接口获取token,再携带token建立WebSocket连接。我说的更稳妥的做法是两者都做——HTTP登录先确认身份,WebSocket连接时再校验token是否有效。注意不要在前端URL后面拼token,那样token很容易被Nginx日志、浏览器历史记录甚至CDN链路记录下来,几乎等于明文泄露。
WebSocket连接建立后的权限控制也是需要提前设计的。比如用户A拉一个群,如果直接向群里任何人mention或者调用私聊接口,必须验证A确实是该会话的成员。我在MessageMapping方法里都调用了conversationMemberMapper.checkMembership,如果返回空直接抛出异常拒绝处理。这个校验在第一个版本就能写,不要拖到上线后补。
6.3 连接数预估和线程池调优
WebSocket是长连接,不像HTTP那样按请求取消释放。单个Tomcat的线程模型在大量WebSocket连接下要重点考虑内存和线程占用。SpringBoot内嵌Tomcat时,每个WebSocket连接会占用一个NIO连接,但Tomcat的NIO线程池默认值是200。这200个线程是所有HTTP和WebSocket共享的,如果同时有大量HTTP请求和大量在线WebSocket连接,就会出现可用线程不足导致连接超时。我们压测后把Tomcat的maxThreads调到500,但注意线程数也不是越大越好,因为每个线程都有栈空间,500个线程对内存的占用已经不小了,要根据机器规格动态调整。
更精细化的控制还有Servlet容器的异步请求超时设置。WebSocket的升级请求本身是异步处理,如果异步超时设置过短可能导致数据还没传输完就被断开。这些参数建议在压测阶段就充分暴露,而不是上线后让用户反馈。长连接服务必须做容量评估,这个经验传递给初次搭IM的团队会非常有价值。
6.4 发版与调试的杀手锏
服务端发版升级时,正在连接的WebSocket会话会被强制断开。如果客户端有自动重连机制,回连会瞬间产生高负载。我后来养成一个习惯:发版前先通过Redis标记服务“停服维护”,前端检测到连接断开后弹出友好提示,而不是疯狂自动重连轰炸服务端。维护结束后清掉标记,用户点重连按钮即可。这就是优雅停机和主动降级的思路,在IM这类高实时性的服务里必不可少。
调试时可以启用WebSocket的底层日志。SpringBoot里设置logging.level.org.springframework.web.socket=DEBUG,前端在stompjs里设置client.debug = (msg) => console.log(msg)。能看到STOMP帧的内容对排查问题极其有帮助,比如某个消息是否真的发到了/topic/chat/1001这个地址。线上环境不要打开debug日志,日志量太大了,但开发环境开着debug信息写代码能看透整个消息流转路径。
最终经验总结
如果你只是做Demo或内部工具,这套方案已经绰绰有余。真要支撑上万人的企业IM,还需要引入更完整的消息可靠性机制,比如端到端加密、多端同步、消息搜索引擎等,但核心的通信骨架已经足够稳了。
我做这个项目时最大的感受是,IM系统的复杂度不是高在通信协议本身,而是高在产品需求和技术方案的夹缝里。你以为在写WebSocket,其实在写会话列表的SQL优化;你以为在调心跳参数,其实在调Nginx超时时间。这种全链路的问题沉淀比任何八股文都值钱。
最后一个小技巧送给大家:STOMP的广播地址和订阅地址之间一定要保持严格一致,前端订阅/topic/chat/1001,后端广播/topic/chat/1002,看起来只差一个数字,定位起来可能要花一整天。建议前后端把destination地址统一维护在一个常量文件里,每次修改同步进行。我对这套方案的最大信心来自它的简洁和通用——技术选型没有追新求异,每个环节都是成熟方案的组合。如果你也想在企业系统里加IM能力,从这套方案起步,应该能少走不少弯路。