ARTICLE DETAIL

资讯详情

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

从零搭建即时通讯后端:WebSocket长连接与消息可靠性实践

从零搭建即时通讯后端:WebSocket长连接与消息可靠性实践 好的收到你的需求。我们直接切入正题。在业务迭代中我们常常会遇到一些名字听起来很“随意”的内部项目比如“猫箱你想干嘛”。但剥开外壳这类项目往往对应着一个非常经典且硬核的技术领域即时通讯IM软件的后端开发。本文就以“猫箱”这个内部项目代号为背景完整拆解一个聊天软件从零到可用的核心后端技术方案。本文将围绕WebSocket 长连接管理、消息可靠性投递、在线状态维护、离线消息补拉这四个即时通讯IM系统最核心的工程问题展开提供完整的代码示例、配置方案和线上避坑指南。无论你是学生想入门 IM 开发还是后端工程师需要在业务中快速落地一个聊天功能本文都能提供一套闭环的实战参考。1. “聊天软件”背后的核心技术挑战先来对齐一下概念。所谓“聊天软件分享”如果从产品经理嘴里说出来可能就是“做一个像微信一样的聊天框”。但如果从后端工程师的视角来看这句话翻译过来其实是客户端 A 发了一条消息服务端如何以毫秒级的延迟把消息推送给客户端 B如果客户端 B 当时不在线比如 App 被杀死消息存哪里什么时候补发如果网络不稳定消息发到一半断了怎么保证消息不丢且不重复服务器怎么知道那么多用户里谁在线、谁离线这几个问题几乎涵盖了 IM 后端 80% 的工作量。很多刚接触 IM 开发的同学第一反应是“用 HTTP 轮询不就行了吗”但实际生产环境中HTTP 轮询的实时性差、服务端压力大而且无法实现服务端主动推送。因此业界的主流方案是使用WebSocket作为消息推送通道配合Redis做在线状态存储配合MySQL或 MongoDB做消息持久化。为了便于理解和后续实操我们暂且把这个项目代号定为“猫箱 IM”。本文的所有代码示例都是为了解决“猫箱 IM”在开发中遇到的真实问题而设计的。2. 环境准备与项目架构设计2.1 技术栈选型说明在动手写代码之前我们需要明确技术栈。本文的示例项目采用以下环境读者在本地复现时请根据实际安装版本微调核心逻辑不受版本影响。开发语言Java 8示例代码使用 Java 11 语法核心框架Spring Boot 2.7.x实时通信WebSocketSpring 自带支持数据存储MySQL 8.x消息记录持久化、Redis 5.x在线状态、会话管理构建工具Maven 3.6IDEIntelliJ IDEA 或 Eclipse2.2 整体架构分层我们先看一张简化的逻辑架构图用它来理解消息的流转路径。这里不使用 Mermaid我们用文字和列表来拆解。整体上一个消息从发送端到接收端会经历四个核心环节接入层WebSocket 服务端负责维护与客户端的 TCP 长连接接收客户端上报的事件上线、下线、心跳以及下行推送消息。逻辑层消息处理中心收到消息后先做合法性校验是否好友、是否被拉黑、内容过滤敏感词、然后生成全局唯一的msgId。存储层持久化将消息内容异步写入 MySQL 或消息队列确保离线用户可以拉取。这里我们为了简化采用同步写 MySQL 的方式生产环境可引入 RocketMQ / RabbitMQ 做解耦。推送路由层状态查询根据接收者的userId查询 Redis判断接收者当前连接在哪一台服务器节点上长连接是四层负载均衡无法跨节点直接推送然后将消息转发到对应节点由该节点推送给客户端。这个流程可以用下面这个有序列表来表示步骤 1用户 A 通过 WebSocket 发送一条 JSON 格式的消息到服务器。步骤 2服务器校验消息格式生成msgId将消息内容写入 MySQL。步骤 3服务器查询 Redis 中用户 B 的在线状态和所在节点。步骤 4服务器将消息通过内部 RPC 或 Redis 发布订阅Pub/Sub转发到用户 B 所在节点。步骤 5用户 B 所在节点的 WebSocket 连接将消息推送给客户端 B。理解了这个流程接下来我们开始搭建 Spring Boot 项目环境。3. 核心依赖与基础配置3.1 创建 Spring Boot 项目并添加依赖在pom.xml中添加 WebSocket 和 Redis 的相关依赖。这里我们使用 Spring Boot 官方 starter 简化配置!-- 文件路径pom.xml -- dependencies !-- Spring Boot Web 基础依赖 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency !-- WebSocket 依赖 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-websocket/artifactId /dependency !-- Redis 依赖用于状态管理 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency !-- JSON 处理工具 -- dependency groupIdcom.alibaba/groupId artifactIdfastjson/artifactId version1.2.83/version /dependency /dependencies这里需要特别说明一点spring-boot-starter-websocket是 Spring 对 WebSocket 协议的封装。它底层在 Servlet 容器中运行利用ServerEndpoint注解或TextWebSocketHandler来处理连接事件。我们下面采用的是TextWebSocketHandler方式这种方式更利于与 Spring IOC 容器整合方便注入 Service 层组件。3.2 配置文件 application.yml在src/main/resources/application.yml中写入以下基础配置。配置内容无需过于复杂重点是为 WebSocket 连接准备端口和 Redis 连接信息。# 文件路径src/main/resources/application.yml server: # WebSocket 服务占用端口 port: 8080 spring: application: name: catbox-im-server redis: # Redis 服务地址按实际环境修改 host: 127.0.0.1 port: 6379 # password: 你的密码如果没有密码就注释掉 database: 0 timeout: 3000ms lettuce: pool: max-active: 8 max-idle: 8 min-idle: 0 # 自定义配置WebSocket 连接路径 catbox: ws: # 前端连接时的端点路径 endpoint: /ws/chat3.3 配置 WebSocket 端点Spring Boot 集成 WebSocket 的核心是注册一个WebSocketHandler并指定它处理的 URL 路径。我们新建一个WebSocketConfig配置类用来注册我们后续会定义的ChatWebSocketHandler。// 文件路径src/main/java/com/catbox/im/config/WebSocketConfig.java package com.catbox.im.config; import com.catbox.im.handler.ChatWebSocketHandler; import org.springframework.context.annotation.Configuration; import org.springframework.web.socket.config.annotation.EnableWebSocket; import org.springframework.web.socket.config.annotation.WebSocketConfigurer; import org.springframework.web.socket.config.annotation.WebSocketHandlerRegistry; import javax.annotation.Resource; /** * WebSocket 配置类 * 作用将自定义的 ChatWebSocketHandler 注册到指定 URL 路径上 */ Configuration EnableWebSocket public class WebSocketConfig implements WebSocketConfigurer { Resource private ChatWebSocketHandler chatWebSocketHandler; Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { registry.addHandler(chatWebSocketHandler, /ws/chat) // 允许跨域访问开发环境配置为 true生产环境请设置为具体的域名白名单 .setAllowedOrigins(*); } }上面的配置中setAllowedOrigins(*)在开发环境很方便但生产环境建议关闭跨域或指定可信来源避免被恶意页面发起 WebSocket 连接。4. 核心代码实现从消息模型到连接管理框架搭好了接下来是重头戏编写消息处理的完整逻辑。我们将按照“消息数据结构 - 在线会话管理 - WebSocket 处理器 - 消息推送服务”的顺序一步步实现一个可运行的单聊功能。4.1 定义统一消息协议前后端交互必须约定一个统一的消息格式。我们定义一个MessageDTO用于承载客户端发送的消息内容和服务端下发的消息内容。// 文件路径src/main/java/com/catbox/im/model/MessageDTO.java package com.catbox.im.model; import java.io.Serializable; /** * 即时通讯统一消息体 */ public class MessageDTO implements Serializable { private static final long serialVersionUID 1L; /** * 消息类型 * 0 - 心跳消息 * 1 - 文本消息 * 2 - 图片消息 * 3 - 上线通知 * 4 - 离线通知 */ private Integer type; /** * 发送者用户ID */ private Long senderId; /** * 接收者用户ID */ private Long receiverId; /** * 全局唯一消息ID由服务端生成 */ private String msgId; /** * 消息内容 */ private String content; /** * 消息发送时间戳毫秒 */ private Long timestamp; public Integer getType() { return type; } public void setType(Integer type) { this.type type; } public Long getSenderId() { return senderId; } public void setSenderId(Long senderId) { this.senderId senderId; } public Long getReceiverId() { return receiverId; } public void setReceiverId(Long receiverId) { this.receiverId receiverId; } public String getMsgId() { return msgId; } public void setMsgId(String msgId) { this.msgId msgId; } public String getContent() { return content; } public void setContent(String content) { this.content content; } public Long getTimestamp() { return timestamp; } public void setTimestamp(Long timestamp) { this.timestamp timestamp; } Override public String toString() { return MessageDTO{ type type , senderId senderId , receiverId receiverId , msgId msgId \ , content content \ , timestamp timestamp }; } }4.2 在线连接管理WebSocketSession 持有与路由既然服务端要主动推送消息就必须在内存中保存每个客户端对应的WebSocketSession对象。单机环境下我们可以直接用ConcurrentHashMap管理。但在生产环境一台服务器撑不下所有用户因此我们需要把“用户ID”和“服务器节点信息”的映射关系存到 Redis 中这样消息入口服务器才能知道该把消息转发给谁。我们先实现一个SessionManager组件它的职责有两个维护当前服务器节点内部的用户连接映射。定时将节点的在线用户信息上报到 Redis。// 文件路径src/main/java/com/catbox/im/config/SessionManager.java package com.catbox.im.config; import org.springframework.stereotype.Component; import org.springframework.web.socket.WebSocketSession; import java.io.IOException; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; /** * 会话管理器负责当前服务器节点上的 WebSocket 连接维护 */ Component public class SessionManager { /** * 当前节点在线用户连接池 * key用户ID * value用户的 WebSocket 会话 */ private static final MapLong, WebSocketSession ONLINE_SESSIONS new ConcurrentHashMap(); /** * 添加会话 */ public void addSession(Long userId, WebSocketSession session) { ONLINE_SESSIONS.put(userId, session); } /** * 移除会话 */ public void removeSession(Long userId) { ONLINE_SESSIONS.remove(userId); } /** * 获取会话 */ public WebSocketSession getSession(Long userId) { return ONLINE_SESSIONS.get(userId); } /** * 判断用户是否在当前节点在线 */ public boolean isOnline(Long userId) { WebSocketSession session ONLINE_SESSIONS.get(userId); return session ! null session.isOpen(); } /** * 获取当前节点连接数量 */ public int getOnlineCount() { return ONLINE_SESSIONS.size(); } /** * 关闭用户连接 */ public void closeSession(Long userId) throws IOException { WebSocketSession session ONLINE_SESSIONS.remove(userId); if (session ! null session.isOpen()) { session.close(); } } }4.3 核心消息处理器ChatWebSocketHandler接下来是 IM 服务器最重要的一个类ChatWebSocketHandler。它负责处理连接建立、消息接收、连接断开等生命周期事件。在handleTextMessage方法中我们接收到一条消息后的处理步骤是将 JSON 字符串解析为MessageDTO。如果消息类型是心跳type0则返回一个心跳 ack保持连接活跃。如果消息类型是上线通知type3则将该用户加入 SessionManager并更新 Redis 中的在线状态。如果消息类型是文本消息type1则调用MessagePushService完成消息持久化和实时推送。// 文件路径src/main/java/com/catbox/im/handler/ChatWebSocketHandler.java package com.catbox.im.handler; import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSONObject; import com.catbox.im.config.SessionManager; import com.catbox.im.model.MessageDTO; import com.catbox.im.service.MessagePushService; import org.springframework.stereotype.Component; import org.springframework.web.socket.CloseStatus; import org.springframework.web.socket.TextMessage; import org.springframework.web.socket.WebSocketSession; import org.springframework.web.socket.handler.TextWebSocketHandler; import javax.annotation.Resource; /** * WebSocket 消息处理器 * 职责处理连接生命周期事件与上行消息 */ Component public class ChatWebSocketHandler extends TextWebSocketHandler { Resource private SessionManager sessionManager; Resource private MessagePushService messagePushService; /** * 连接建立后触发 * 解析 URL 参数中的 userId作为连接的身份标识 */ Override public void afterConnectionEstablished(WebSocketSession session) throws Exception { // 握手时通过查询参数传递用户IDws://localhost:8080/ws/chat?userId1001 String query session.getUri().getQuery(); Long userId parseUserIdFromQuery(query); if (userId null) { session.close(CloseStatus.BAD_DATA); return; } // 将登录用户加入当前节点的连接池 sessionManager.addSession(userId, session); // 模拟打印日志生产环境请使用 Slf4j System.out.println(用户 userId 已连接当前节点在线人数 sessionManager.getOnlineCount()); } /** * 接收客户端消息 */ Override protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception { // 1. 解析消息为统一对象 MessageDTO messageDTO JSON.parseObject(message.getPayload(), MessageDTO.class); // 2. 根据消息类型分发处理 if (messageDTO.getType() 0) { // 心跳消息原样返回一个 pong JSONObject pong new JSONObject(); pong.put(type, 0); pong.put(msg, pong); session.sendMessage(new TextMessage(pong.toJSONString())); } else if (messageDTO.getType() 1) { // 文本消息交给消息推送服务处理 messagePushService.handleSingleChatMessage(messageDTO); } else if (messageDTO.getType() 3) { // 上线通知更新 Redis 在线状态 messagePushService.processUserOnline(messageDTO.getSenderId()); } } /** * 连接断开后触发 * 清理连接池中的会话并更新 Redis 离线状态 */ Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception { String query session.getUri().getQuery(); Long userId parseUserIdFromQuery(query); if (userId ! null) { sessionManager.removeSession(userId); messagePushService.processUserOffline(userId); System.out.println(用户 userId 已断开连接); } } /** * 从 URL 查询参数中提取 userId * 例如ws://localhost:8080/ws/chat?userId1001 */ private Long parseUserIdFromQuery(String query) { if (query null || query.isEmpty()) { return null; } String[] pairs query.split(); for (String pair : pairs) { String[] kv pair.split(); if (kv.length 2 userId.equals(kv[0])) { return Long.parseLong(kv[1]); } } return null; } }4.4 消息推送与持久化服务消息处理的核心逻辑在MessagePushService中。这里我们实现两个关键方法handleSingleChatMessage(MessageDTO)处理单聊消息先落库再尝试实时推送。sendMessageToUser(Long receiverId, MessageDTO)向指定用户推送消息先查本地再查全局状态。这是 IM 系统的“心脏”我们尽量把逻辑写清楚。// 文件路径src/main/java/com/catbox/im/service/MessagePushService.java package com.catbox.im.service; import com.alibaba.fastjson.JSON; import com.catbox.im.config.SessionManager; import com.catbox.im.model.MessageDTO; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.stereotype.Service; import org.springframework.web.socket.TextMessage; import org.springframework.web.socket.WebSocketSession; import javax.annotation.Resource; import java.util.UUID; /** * 消息推送服务 * 职责消息落库、状态查询、实时推送 */ Service public class MessagePushService { Resource private SessionManager sessionManager; Resource private StringRedisTemplate stringRedisTemplate; /** * 处理单聊消息 */ public void handleSingleChatMessage(MessageDTO messageDTO) { // 1. 为消息生成全局唯一 ID String msgId UUID.randomUUID().toString().replace(-, ); messageDTO.setMsgId(msgId); // 2. 防止前端伪造时间戳 if (messageDTO.getTimestamp() null) { messageDTO.setTimestamp(System.currentTimeMillis()); } // 3. 消息落库此处为核心片段实际项目中通常为异步写入 MySQL saveMessageToDB(messageDTO); // 4. 尝试实时推送 boolean pushed sendMessageToUser(messageDTO.getReceiverId(), messageDTO); // 5. 如果用户不在线可以记录一条离线日志等待用户上线后重新拉取 if (!pushed) { // 生产环境这里可以发送到消息队列等用户上线后消费 System.out.println(用户 messageDTO.getReceiverId() 当前不在线消息已存储待上线补拉); } } /** * 向指定用户推送消息 * 返回 true 表示推送成功false 表示用户不在线 */ public boolean sendMessageToUser(Long receiverId, MessageDTO messageDTO) { // 1. 先查当前节点的连接池 WebSocketSession session sessionManager.getSession(receiverId); if (session ! null session.isOpen()) { try { session.sendMessage(new TextMessage(JSON.toJSONString(messageDTO))); return true; } catch (Exception e) { e.printStackTrace(); return false; } } // 2. 如果当前节点没有该用户查 Redis 中的全局在线状态 // 这里为了演示简单我们直接返回 false。 // 在实际分布式环境中需要根据 Redis 中记录的节点ID将消息转发给对应的服务器节点。 String redisKey im:online:user: receiverId; Boolean isOnline stringRedisTemplate.hasKey(redisKey); if (Boolean.TRUE.equals(isOnline)) { // 理论上是分布式转发逻辑示例代码不再展开 System.out.println(用户在线但连接在其他节点需要做节点间转发); // 伪代码rpcClient.sendToNode(nodeId, messageDTO); } return false; } /** * 用户上线更新 Redis 状态 */ public void processUserOnline(Long userId) { String redisKey im:online:user: userId; // 记录在线状态并设置过期时间防止异常掉线后 key 残留 stringRedisTemplate.opsForValue().set(redisKey, String.valueOf(System.currentTimeMillis()), 30, java.util.concurrent.TimeUnit.MINUTES); System.out.println(用户 userId 上线Redis 状态已更新); } /** * 用户下线删除 Redis 状态 */ public void processUserOffline(Long userId) { String redisKey im:online:user: userId; stringRedisTemplate.delete(redisKey); System.out.println(用户 userId 下线Redis 状态已删除); } /** * 模拟消息落库 * 真实场景中请使用 MyBatis-Plus 或 JPA 写入 MySQL 的 message 表 */ private void saveMessageToDB(MessageDTO messageDTO) { // 这里只是模拟实际项目中需要调用 Mapper 层。 // 示例 SQL 见文章后续章节的建表语句。 System.out.println(消息落库成功 messageDTO.getMsgId()); } }5. 离线消息与消息可靠送达聊完了在线实时推送我们来看看 IM 系统中另一个无法回避的场景当接收方不在线时消息如何补发5.1 离线消息存储表设计在 MySQL 中我们可以设计一张简单的im_message表存储所有聊天记录。当用户上线时前端会传入一个lastMsgId本地最后一条消息的 ID或时间戳服务端根据这个游标去数据库查询增量消息再批量推送给客户端。以下是消息表的建表语句作为参考。这条语句只是一个简单的思路演示生产环境需要考虑分库分表。-- 文件路径docs/sql/im_message.sql CREATE TABLE im_message ( id bigint(20) NOT NULL AUTO_INCREMENT COMMENT 自增主键, msg_id varchar(64) NOT NULL COMMENT 全局唯一消息ID, sender_id bigint(20) NOT NULL COMMENT 发送者用户ID, receiver_id bigint(20) NOT NULL COMMENT 接收者用户ID, content text COMMENT 消息内容, message_type tinyint(4) DEFAULT 1 COMMENT 消息类型1-文本2-图片3-语音, status tinyint(4) DEFAULT 0 COMMENT 消息状态0-未读1-已读, create_time datetime DEFAULT CURRENT_TIMESTAMP COMMENT 创建时间, PRIMARY KEY (id), UNIQUE KEY uk_msg_id (msg_id), KEY idx_receiver_time (receiver_id, create_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT聊天消息记录表;5.2 离线消息补拉核心逻辑当用户上线建立 WebSocket 连接之后客户端会主动发送一个“同步离线消息”的请求。这里我们简化处理通过一个 HTTP 接口或 WebSocket 消息类型来实现。核心的查询逻辑是根据接收者 ID 和本地最大消息 ID查找数据库中比这个 ID 大的所有消息。// 这是 Service 层方法的核心片段 // 文件路径src/main/java/com/catbox/im/service/MessageSyncService.java片段 public ListMessageDTO pullOfflineMessages(Long receiverId, Long lastMsgId, Integer limit) { // 请根据项目实际情况注入 Mapper // 这里演示 MyBatis-Plus 的 LambdaQueryWrapper 写法 // LambdaQueryWrapperMessageEntity wrapper new LambdaQueryWrapper(); // wrapper.eq(MessageEntity::getReceiverId, receiverId) // .gt(MessageEntity::getId, lastMsgId) // .orderByAsc(MessageEntity::getId) // .last(limit limit); // ListMessageEntity entities messageMapper.selectList(wrapper); // 将实体列表转换成 MessageDTO 列表返回 return new ArrayList(); }5.3 消息去重策略为了保证消息不丢很多客户端逻辑会在断线重连后向上取一段时间的消息拉取。例如本地已有 1~100 条消息断线期间错过了 101~105 条重连后客户端会传lastMsgId100来补拉。但如果是推送超时导致的消息重复客户端需要根据msgId做去重。因此我们在设计MessageDTO时特意要求服务端生成全局唯一的msgId客户端在展示消息时需要用msgId作为本地缓存的 key而不是用自增 ID 或时间戳。6. 常见问题与线上排障锦囊实战中IM 系统最容易遇到的问题我整理成了一张排查表方便大家直接对照解决。问题现象常见原因解决思路WebSocket 连接一直处于CONNECTING状态前端连接 URL 错误服务端端口未开放路径注册错误。检查ws://还是wss://协议检查路径是否与registry.addHandler中一致用浏览器控制台网络面板看握手请求是否被拒。连接一建立就被服务端关闭afterConnectionEstablished中session.close(CloseStatus.BAD_DATA)被触发通常是userId参数缺失。检查 WebSocket 连接 URL 是否拼接了?userIdxxx检查参数解析逻辑是否区分了多个参数。消息只能发给同一台服务器上的用户多节点部署时没有实现节点间消息转发只在本地SessionManager中查找了会话。引入 Redis Pub/Sub 或消息队列将所有节点的在线状态和消息路由集中管理将消息按receiverId哈希转发到对应节点。心跳包发送后连接仍然断开服务端没有处理心跳消息或者心跳间隔大于 Nginx/负载均衡器的空闲超时时间。在handleTextMessage中识别心跳类型并返回 ack合理设置心跳间隔建议 30s~60s。Redis 中在线用户越来越多出现脏数据客户端异常断网例如弱网掉线服务端没有及时执行afterConnectionClosed。给 Redis key 设置过期时间如 30 分钟增加服务端心跳检测任务定期清理失效连接。补拉离线消息时消息重复客户端没有用msgId做幂等去重。服务端下发消息时保证msgId全局唯一客户端根据msgId做去重展示。7. IM 工程化最佳实践与生产建议把功能跑通只是第一步。真正上了生产环境以下几个问题会让你少走很多弯路。7.1 关于 WebSocket 连接鉴权很多初学者容易忽略WebSocket 握手本质上是一个 HTTP GET 请求。因此我们可以利用这个特点在握手阶段完成用户鉴权。常见的做法是客户端在 URL 上携带一个短期有效的 Token服务端在HandshakeInterceptor中校验 Token如果校验失败直接拒绝握手。这样可以避免把长期有效的用户密码暴露在 URL 中。这里给出一个HandshakeInterceptor的配置思路不展开完整代码。// 核心思路实现 HandshakeInterceptor在 beforeHandshake 中校验 Token public class AuthHandshakeInterceptor implements HandshakeInterceptor { Override public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, MapString, Object attributes) throws Exception { // 从 URL 中获取 token String token request.getURI().getQuery(); // 校验 token 合法性例如调用 Redis 查询或 JWT 解析 // 校验通过后将 userId 放入 attributes供 WebSocketSession 使用 // 校验失败返回 false拒绝握手 return true; } }7.2 关于心跳机制与连接保活生产环境几乎必然存在 Nginx、SLB 等负载均衡组件。这些组件如果在一定时间内没有收到流量会主动断开空闲连接。因此客户端必须实现应用层心跳。我们的建议是客户端每 30 秒发送一个心跳包。服务端收到心跳后返回一个 pong 消息。服务端如果 90 秒内没有收到任何消息含心跳则主动关闭该连接。客户端收到连接关闭事件后执行指数退避重连。7.3 关于消息存储的异步化在高并发场景下不要在 WebSocket 的处理线程里同步写 MySQL。因为数据库写入涉及磁盘 IO耗时不可控会拖慢整个消息链路的延迟。正确做法是WebSocket 处理器将消息写入 Redis List 或直接发送到消息队列Kafka/RabbitMQ。消费者异步消费消息批量写入 MySQL。如果写入失败可以重试保证消息最终一致。7.4 关于日志与监控IM 服务的日志非常重要。建议对以下三个核心指标做监控节点在线连接数通过SessionManager.getOnlineCount()暴露给监控系统如 Prometheus。WebSocket 握手成功率用来评估鉴权逻辑是否正常。消息推送延迟在消息入口打一个时间戳在下发时计算延迟超出阈值报警。这一点往往被业务快速迭代所忽略但一旦线上出问题没有监控日志的 IM 系统排查起来就像大海捞针。7.5 生产环境部署架构参考最后给出一个适合中小团队上线的部署建议使用 Nginx 作为 WebSocket 反向代理配置proxy_set_header Upgrade $http_upgrade;和proxy_set_header Connection upgrade;。Nginx 负载均衡算法推荐ip_hash这样可以尽量保证同一个用户的请求落在同一台服务器节点上降低跨节点转发的概率。如果用户规模较大建议单独申请 WebSocket 专用域名与业务 API 域名隔离方便做限流和安全策略。实际上“猫箱 IM”这类项目的技术难点并不在于某个单独的 API 怎么调用而在于你对连接生命周期和消息一致性这两个核心问题的理解深度。掌握了本文中的 Session 管理、消息落库、离线补拉这三个模块你已经具备开发一个生产级聊天软件后端的雏形能力。如果这篇文章帮你理清了思路或者解决了你的实际问题可以收藏备用。后续也可以继续深入 Netty 高性能网络编程、消息队列削峰填谷以及分布式 IM 架构设计等方向这些都是聊天软件后端进阶的必经之路。
返回列表