你是不是觉得站内信推送系统很简单?不就是用户发个消息,系统存一下,然后推给另一个用户吗?很多开发者一开始都这么想,直到真正动手时才发现处处是坑:消息延迟、已读状态不同步、海量数据下的性能瓶颈、推送失败如何补偿……一个看似简单的功能,背后却涉及系统架构、数据一致性、实时性和扩展性等多个维度的设计挑战。
这篇文章要解决的,正是这个“看起来简单,做起来头疼”的问题。我们将从一个真实的业务场景出发,彻底拆解一个高可用、可扩展的站内信推送系统该如何设计与实现。读完本文,你将不仅知道如何用代码实现收发消息,更能掌握一套应对高并发、保证消息必达、支持多端同步的完整设计思路与工程实践。无论你是要为一个快速发展的社区、一个电商客服系统,还是一个内部协作工具搭建消息通道,这里都有你需要的答案。
1. 站内信系统:远不止“发消息”那么简单
在深入代码之前,我们必须先厘清“站内信推送系统”的核心边界与设计目标。它不是一个简单的INSERT加SELECT操作。一个成熟的生产级系统,需要同时满足以下几个看似矛盾的需求:
- 高实时性:用户发送消息后,接收方应近乎实时地感知。
- 高可靠性:消息不能丢失,必须保证“至少送达一次”(At-Least-Once Delivery)。
- 状态一致性:消息的“已读/未读”状态必须在所有客户端(Web、App)间实时同步。
- 海量数据支撑:系统需要能平滑应对用户量和消息量的指数级增长。
- 低延迟与高并发:在万人同时在线聊天的场景下,系统不能雪崩。
传统的、基于数据库轮询(Polling)的简单方案(例如,前端每5秒查询一次数据库“是否有新消息”)在以上任何一点面前都会迅速崩溃。它会给数据库带来巨大压力,实时性差,且无法有效同步状态。
因此,现代站内信系统的核心设计范式已经转向了“事件驱动”和“长连接推送”。简单来说,系统的工作流变成了这样:发送方触发一个“发送消息”事件 -> 系统持久化消息并生成一个“新消息”事件 -> 通过长连接通道实时推送给在线的接收方。对于离线的接收方,则在其下次上线时主动拉取未读消息。
接下来,我们将从概念到实现,一步步构建这个系统。
2. 核心概念与架构设计
2.1 核心数据模型
首先定义最核心的实体:消息(Message)。
-- 文件路径:/sql/create_table_messages.sql CREATE TABLE `user_message` ( `id` bigint(20) NOT NULL AUTO_INCREMENT COMMENT '消息ID,主键', `sender_id` bigint(20) NOT NULL COMMENT '发送者用户ID', `receiver_id` bigint(20) NOT NULL COMMENT '接收者用户ID', `content` text NOT NULL COMMENT '消息内容(可存储JSON或纯文本)', `content_type` tinyint(4) NOT NULL DEFAULT '1' COMMENT '消息类型:1-文本,2-图片,3-文件...', `conversation_id` varchar(128) NOT NULL COMMENT '会话ID,用于标识两个用户间的唯一对话通道,通常为 sorted(sender_id, receiver_id)', `status` tinyint(4) NOT NULL DEFAULT '0' COMMENT '消息状态:0-发送中,1-已送达,2-已读', `read_at` datetime DEFAULT NULL COMMENT '阅读时间', `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间', `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间', PRIMARY KEY (`id`), KEY `idx_conversation_id` (`conversation_id`), KEY `idx_receiver_status` (`receiver_id`, `status`), KEY `idx_created_at` (`created_at`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='用户消息表';关键设计解析:
conversation_id:这是优化查询的关键。将“查询用户A和用户B的所有消息”这个复杂条件((sender=A AND receiver=B) OR (sender=B AND receiver=A)),简化为对conversation_id的单列查询,性能提升巨大。status:明确的消息状态机(发送中->已送达->已读)是实现可靠性追踪的基础。- 索引设计:
idx_receiver_status索引专门用于高效查询某个用户的未读消息。idx_created_at用于消息分页拉取。
2.2 系统架构总览
一个典型的、解耦的站内信推送系统架构如下所示:
[客户端 App/Web] | | (1. 建立长连接) | [WebSocket/Grpc-Web 网关层] —— 负责维护海量用户长连接,管理会话。 | | (2. 路由消息事件) | [消息推送服务] —————— (3. 持久化消息) —————— [MySQL/PostgreSQL 消息存储] | | | (4. 发布新消息事件) | (5. 离线消息拉取) | | [消息队列 (如 RabbitMQ/Kafka)] | | | | (6. 消费事件,查找在线连接) | | | [连接管理服务] ——— (7. 通过网关推送) ———→ [WebSocket/Grpc-Web 网关层] | | (8. 推送到客户端) | [客户端 App/Web]各组件职责:
- 网关层:技术选型可以是 Netty 实现的 WebSocket 服务器,或 Spring Boot 集成的
@ServerEndpoint,亦或是专门的 Go 语言网关。它负责最底层的连接保持、心跳检测和帧解析。 - 消息推送服务:业务逻辑的核心。接收发送请求,处理消息,写入数据库,并向消息队列发布事件。
- 消息队列:系统的“中枢神经”。它解耦了消息的“生产”(写入)和“消费”(推送),使得在线推送和离线存储可以异步处理,提高了系统的吞吐量和抗压能力。
- 连接管理服务:维护一个“用户ID -> 连接实例”的映射关系(通常存在于 Redis 中),当需要推送时,能快速找到接收方当前所在的网关节点和连接。
这套架构的核心优势在于解耦和异步。发送消息的API可以快速响应,而耗时的推送任务由下游消费者异步完成。即使推送服务暂时不可用,消息也已安全持久化,不会丢失。
3. 环境准备与核心技术栈
在开始编码前,我们需要搭建基础环境。本文将以Spring Boot作为后端框架,Netty实现 WebSocket 网关,RabbitMQ作为消息队列,Redis用于连接管理和会话缓存,MySQL作为主存储。
3.1 依赖清单 (Maven pom.xml)
<!-- 文件路径:pom.xml --> <dependencies> <!-- Spring Boot 基础 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-websocket</artifactId> </dependency> <!-- 数据持久化 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-jpa</artifactId> </dependency> <dependency> <groupId>mysql</groupId> <artifactId>mysql-connector-java</artifactId> <scope>runtime</scope> </dependency> <!-- 缓存与消息队列 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-redis</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> <!-- Netty (用于更底层的WebSocket实现,可选) --> <dependency> <groupId>io.netty</groupId> <artifactId>netty-all</artifactId> <version>4.1.108.Final</version> </dependency> <!-- 工具类 --> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <optional>true</optional> </dependency> <dependency> <groupId>com.alibaba.fastjson2</groupId> <artifactId>fastjson2</artifactId> <version>2.0.51</version> </dependency> </dependencies>3.2 基础配置 (application.yml)
# 文件路径:src/main/resources/application.yml spring: datasource: url: jdbc:mysql://localhost:3306/message_db?useUnicode=true&characterEncoding=utf8&useSSL=false&serverTimezone=Asia/Shanghai username: root password: your_password driver-class-name: com.mysql.cj.jdbc.Driver jpa: hibernate: ddl-auto: update # 生产环境请改为 validate 或 none,并使用Flyway/Liquibase show-sql: true properties: hibernate: format_sql: true redis: host: localhost port: 6379 password: # 如果有密码则填写 database: 0 lettuce: pool: max-active: 8 max-wait: -1ms max-idle: 8 min-idle: 0 rabbitmq: host: localhost port: 5672 username: guest password: guest listener: simple: acknowledge-mode: manual # 手动ACK,保证消息可靠消费 # 自定义配置 app: websocket: port: 8081 # Netty WebSocket服务器端口4. 核心流程拆解与实现
我们将核心流程拆解为四个关键步骤:建立连接、发送消息、消息持久化与事件发布、实时推送。
4.1 第一步:建立与管理长连接
我们使用 Netty 实现一个轻量级的 WebSocket 服务器,用于维持客户端长连接。
// 文件路径:src/main/java/com/example/message/gateway/WebSocketServer.java @Component public class WebSocketServer { private final EventLoopGroup bossGroup = new NioEventLoopGroup(); private final EventLoopGroup workerGroup = new NioEventLoopGroup(); @Value("${app.websocket.port}") private int port; @PostConstruct public void start() throws InterruptedException { ServerBootstrap bootstrap = new ServerBootstrap(); bootstrap.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .childHandler(new ChannelInitializer<SocketChannel>() { @Override protected void initChannel(SocketChannel ch) { ChannelPipeline pipeline = ch.pipeline(); // 处理HTTP请求和WebSocket握手 pipeline.addLast(new HttpServerCodec()); pipeline.addLast(new HttpObjectAggregator(65536)); pipeline.addLast(new WebSocketServerProtocolHandler("/ws")); // 自定义消息处理器 pipeline.addLast(new WebSocketFrameHandler()); } }); ChannelFuture future = bootstrap.bind(port).sync(); future.channel().closeFuture().addListener(f -> { bossGroup.shutdownGracefully(); workerGroup.shutdownGracefully(); }); } }// 文件路径:src/main/java/com/example/message/gateway/WebSocketFrameHandler.java public class WebSocketFrameHandler extends SimpleChannelInboundHandler<TextWebSocketFrame> { @Override public void channelActive(ChannelHandlerContext ctx) { // 连接建立,可以在此进行认证(例如通过URL参数传递token) String token = getTokenFromUri(ctx.channel()); Long userId = authService.validateToken(token); if (userId != null) { // 将 userId 与 Channel 绑定 ChannelSupervise.addChannel(userId, ctx.channel()); // 将映射关系存入Redis: user:online:{userId} -> gatewayNodeId:channelId redisTemplate.opsForValue().set("user:online:" + userId, ctx.channel().id().asLongText(), 5, TimeUnit.MINUTES); } else { ctx.close(); } } @Override protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame frame) { // 处理客户端发来的消息,例如心跳包、客户端ACK等 String text = frame.text(); // 解析协议,例如:{"type":"heartbeat"} 或 {"type":"ack", "msgId":"123"} handleClientMessage(ctx, text); } @Override public void channelInactive(ChannelHandlerContext ctx) { // 连接断开,清理资源 Long userId = ChannelSupervise.getUserIdByChannel(ctx.channel()); if (userId != null) { ChannelSupervise.removeChannel(ctx.channel()); redisTemplate.delete("user:online:" + userId); } } // ... 其他方法,如异常处理 }关键点:
- 连接认证:在
channelActive时,必须通过 Token 等手段验证用户身份,并将userId与 Netty 的Channel绑定。 - 连接管理:使用一个全局的
ChannelSupervise类(内部可用ConcurrentHashMap)管理在线连接。更重要的是,必须将userId->channel的映射写入 Redis,并设置一个较短的过期时间(如5分钟)。这样,分布式的推送服务才能通过查询 Redis 找到用户连接所在节点。 - 心跳机制:客户端需要定期发送心跳帧,服务端也需要定时检查连接有效性,及时清理僵尸连接。
4.2 第二步:发送消息的API实现
这是业务入口,它负责接收发送请求,执行业务逻辑,并触发后续流程。
// 文件路径:src/main/java/com/example/message/controller/MessageController.java @RestController @RequestMapping("/api/message") public class MessageController { @Autowired private MessageService messageService; @PostMapping("/send") public ApiResponse<SendMessageResult> sendMessage(@RequestBody SendMessageRequest request) { // 1. 参数校验 if (request.getReceiverId() == null || StringUtils.isBlank(request.getContent())) { return ApiResponse.error("参数错误"); } // 2. 调用服务层 SendMessageResult result = messageService.sendMessage(request); return ApiResponse.success(result); } }// 文件路径:src/main/java/com/example/message/service/impl/MessageServiceImpl.java @Service @Slf4j public class MessageServiceImpl implements MessageService { @Autowired private MessageRepository messageRepository; @Autowired private RabbitTemplate rabbitTemplate; @Transactional(rollbackFor = Exception.class) @Override public SendMessageResult sendMessage(SendMessageRequest request) { Long senderId = getCurrentUserId(); // 从安全上下文获取 Long receiverId = request.getReceiverId(); // 1. 构建并保存消息实体 UserMessage message = new UserMessage(); message.setSenderId(senderId); message.setReceiverId(receiverId); message.setContent(request.getContent()); message.setContentType(request.getContentType()); // 生成会话ID: 保证唯一且有序,例如 “小ID:大ID” message.setConversationId(generateConversationId(senderId, receiverId)); message.setStatus(MessageStatus.SENDING.getCode()); // 初始状态为发送中 message = messageRepository.save(message); log.info("消息持久化成功,消息ID: {}", message.getId()); // 2. 构建消息事件,发送到MQ MessageEvent event = new MessageEvent(); event.setMessageId(message.getId()); event.setSenderId(senderId); event.setReceiverId(receiverId); event.setContent(message.getContent()); event.setEventType("NEW_MESSAGE"); rabbitTemplate.convertAndSend("message.exchange", "message.route.key", event); log.info("消息事件已发布到MQ,接收者: {}", receiverId); // 3. 更新消息状态为“已送达”(如果后续推送成功,会再更新为“已读”) // 此处可以先不更新,等推送服务消费成功后回调更新。为简化,我们先更新。 message.setStatus(MessageStatus.DELIVERED.getCode()); messageRepository.save(message); return new SendMessageResult(message.getId(), true); } private String generateConversationId(Long uid1, Long uid2) { long min = Math.min(uid1, uid2); long max = Math.max(uid1, uid2); return min + ":" + max; } }设计解析:
- 事务边界:消息的持久化必须在数据库事务中完成,确保消息不丢失。
- 异步解耦:服务层不直接处理推送逻辑,而是将
MessageEvent发布到 RabbitMQ。这样,sendMessage方法可以快速返回,用户体验好,系统吞吐量高。 - 事件内容:事件中包含了推送所需的最小数据集(
receiverId,content等),避免消费者再去查库,减少延迟。
4.3 第三步:消息队列与事件消费
我们配置 RabbitMQ 的交换机和队列,并编写消费者来监听新消息事件。
// 文件路径:src/main/java/com/example/message/config/RabbitMQConfig.java @Configuration public class RabbitMQConfig { public static final String EXCHANGE_NAME = "message.exchange"; public static final String QUEUE_NAME = "message.push.queue"; public static final String ROUTING_KEY = "message.route.key"; @Bean public DirectExchange messageExchange() { return new DirectExchange(EXCHANGE_NAME, true, false); // 持久化,不自动删除 } @Bean public Queue messageQueue() { return new Queue(QUEUE_NAME, true, false, false); // 持久化队列 } @Bean public Binding binding() { return BindingBuilder.bind(messageQueue()).to(messageExchange()).with(ROUTING_KEY); } }// 文件路径:src/main/java/com/example/message/consumer/MessagePushConsumer.java @Component @Slf4j public class MessagePushConsumer { @Autowired private RedisTemplate<String, String> redisTemplate; @Autowired private WebSocketPushService pushService; @RabbitListener(queues = RabbitMQConfig.QUEUE_NAME) public void handleMessage(MessageEvent event, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) { try { Long receiverId = event.getReceiverId(); log.info("开始处理推送,接收者ID: {}, 消息ID: {}", receiverId, event.getMessageId()); // 1. 检查接收者是否在线 String connectionKey = redisTemplate.opsForValue().get("user:online:" + receiverId); if (StringUtils.isNotBlank(connectionKey)) { // 2. 在线,通过WebSocket推送 boolean pushSuccess = pushService.pushToUser(receiverId, event); if (pushSuccess) { // 推送成功,可以更新消息状态为“已读”(或由客户端ACK确认) log.info("消息实时推送成功, receiverId: {}", receiverId); } else { // 推送失败,可能用户刚好下线,转入离线逻辑 log.warn("实时推送失败,消息转入离线逻辑, receiverId: {}", receiverId); handleOfflineMessage(event); } } else { // 3. 离线,存储到离线消息库(如Redis有序集合或MySQL扩展表) log.info("用户离线,存储离线消息, receiverId: {}", receiverId); handleOfflineMessage(event); } // 4. 手动确认消息,确保可靠性 channel.basicAck(tag, false); } catch (Exception e) { log.error("消息消费失败, event: {}", event, e); // 处理失败,根据策略决定是重试、记录还是丢弃 try { channel.basicNack(tag, false, true); // 重回队列,重试 } catch (IOException ex) { log.error("消息Nack失败", ex); } } } private void handleOfflineMessage(MessageEvent event) { // 将事件存储到以用户ID为key的Redis List或Sorted Set中 String key = "user:offline:msg:" + event.getReceiverId(); redisTemplate.opsForList().rightPush(key, JSON.toJSONString(event)); // 可以设置过期时间,例如7天 redisTemplate.expire(key, 7, TimeUnit.DAYS); } }关键机制:
- 手动ACK:配置
acknowledge-mode: manual并手动调用basicAck,只有业务处理成功后才确认消息,防止消息丢失。 - 失败重试:消费失败时,通过
basicNack将消息重回队列(或进入死信队列),实现自动重试。 - 在线/离线判断:通过查询 Redis 中是否存在
user:online:{userId}键来判断用户在线状态。这是整个推送链路的核心决策点。
4.4 第四步:WebSocket推送服务
推送服务负责根据connectionKey找到具体的 Netty Channel,并发送数据帧。
// 文件路径:src/main/java/com/example/message/service/WebSocketPushService.java @Service public class WebSocketPushService { public boolean pushToUser(Long userId, MessageEvent event) { // 1. 从连接管理器中获取Channel Channel channel = ChannelSupervise.findChannelByUserId(userId); if (channel == null || !channel.isActive()) { // 连接已失效,清理Redis中的记录 redisTemplate.delete("user:online:" + userId); return false; } // 2. 构建推送协议 PushDTO pushDTO = new PushDTO(); pushDTO.setType("NEW_MESSAGE"); pushDTO.setData(event); String pushMessage = JSON.toJSONString(pushDTO); // 3. 通过Netty Channel发送 if (channel.isWritable()) { channel.writeAndFlush(new TextWebSocketFrame(pushMessage)); log.debug("WebSocket推送成功, userId: {}, msgId: {}", userId, event.getMessageId()); return true; } else { log.warn("Channel不可写,推送失败, userId: {}", userId); return false; } } }至此,一个完整的“发送->持久化->事件通知->实时推送”的核心流程已经闭环。
5. 进阶功能与最佳实践
实现基础流程后,一个健壮的系统还需要考虑更多细节。
5.1 消息的可靠性与状态同步
“已送达”和“已读”状态如何准确同步?
方案:客户端ACK机制。
- 服务端推送消息时,携带一个唯一的
msgId。 - 客户端收到后,立即回传一个ACK命令:
{"type": "ack", "msgId": "123"}。 - 服务端的
WebSocketFrameHandler收到ACK后,更新数据库中对应消息的状态为已读,并更新read_at时间。
// 客户端ACK消息处理示例 private void handleClientAck(Long userId, String msgId) { messageRepository.updateStatusToRead(msgId, new Date()); // 可选:广播该消息的已读状态给发送者(用于实现“对方已读”提示) notifySenderMessageRead(msgId); }5.2 离线消息拉取
用户上线后,如何获取离线期间的消息?
方案:上线后主动拉取 + 离线队列。
- 在
WebSocketFrameHandler.channelActive(用户连接建立)时,除了绑定连接,还需要检查该用户是否存在离线消息(即检查Redis中user:offline:msg:{userId}这个List)。 - 如果存在,则一次性或分批次将这些消息推送给用户。
- 推送成功后,从Redis中删除已推送的消息。
// 连接建立后拉取离线消息 private void pushOfflineMessagesAfterLogin(Long userId, Channel channel) { String offlineKey = "user:offline:msg:" + userId; List<String> offlineMsgJsons = redisTemplate.opsForList().range(offlineKey, 0, -1); if (CollectionUtils.isEmpty(offlineMsgJsons)) { return; } for (String msgJson : offlineMsgJsons) { MessageEvent event = JSON.parseObject(msgJson, MessageEvent.class); pushToChannel(channel, event); // 通过当前channel推送 } // 全部推送成功后,清空离线队列 redisTemplate.delete(offlineKey); }5.3 会话列表与未读计数
这是产品层的高频需求。其核心是空间换时间,避免每次打开列表都去聚合查询。
方案:使用Redis维护会话摘要和未读数。
- 会话列表缓存:为每个用户维护一个
Sorted Set,Key为user:conversations:{userId},Score为最后一条消息的时间戳,Value为会话ID。每次收发消息,都更新这个集合。 - 未读计数:为每个用户在每个会话维护一个计数器,Key为
user:unread:{userId}:{conversationId}。收到新消息时INCR,消息被读后SET 0。
// 发送消息时更新会话列表和未读计数 public void updateConversationAndUnread(Long senderId, Long receiverId, String conversationId) { long now = System.currentTimeMillis(); String userConversationKey = "user:conversations:" + receiverId; // 更新接收者的会话列表(按时间排序) redisTemplate.opsForZSet().add(userConversationKey, conversationId, now); // 增加接收者的未读计数 String unreadKey = "user:unread:" + receiverId + ":" + conversationId; redisTemplate.opsForValue().increment(unreadKey); // 可选:限制会话列表长度,只保留最近的50个 redisTemplate.opsForZSet().removeRange(userConversationKey, 0, -51); }5.4 消息历史记录分页查询
查询两人之间的历史消息,利用好conversation_id索引。
// 文件路径:src/main/java/com/example/message/repository/MessageRepository.java @Repository public interface MessageRepository extends JpaRepository<UserMessage, Long> { // 基于游标的分页查询,性能优于 LIMIT offset, size @Query(value = "SELECT * FROM user_message WHERE conversation_id = :conversationId AND id < :lastId ORDER BY id DESC LIMIT :size", nativeQuery = true) List<UserMessage> findHistoryMessages(@Param("conversationId") String conversationId, @Param("lastId") Long lastId, @Param("size") int size); }6. 常见问题与排查思路
在开发和运维过程中,你几乎一定会遇到以下问题:
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| 消息发送成功,但对方收不到 | 1. 接收方长连接已断开。 2. Redis中在线状态丢失。 3. MQ消费者堆积或宕机。 4. 推送服务到网关的网络问题。 | 1. 检查接收方客户端网络和心跳。 2. 查看Redis中 user:online:{userId}键是否存在及过期时间。3. 查看MQ管理界面,检查队列积压情况和消费者状态。 4. 查看推送服务日志,确认是否调用了 pushToUser及返回值。 | 1. 优化心跳机制,客户端实现断线重连。 2. 确保连接建立和断开时,Redis键的设置和删除是原子操作。 3. 增加消费者实例,监控消费延迟。 4. 确保内网服务间网络畅通,考虑服务注册发现。 |
| 消息重复消费 | 1. MQ消息被重复投递(网络问题导致ACK未送达)。 2. 消费者处理超时,触发了重试机制。 | 1. 检查消费者日志,看同一条消息是否被处理多次。 2. 检查业务逻辑是否有幂等性设计。 | 实现消费幂等性:在消费前,先查库判断该messageId的状态是否已处理。或在Redis中设置一个已处理标记msg:processed:{messageId},使用SETNX命令。 |
| 数据库CPU/IO压力高 | 1. 消息表缺乏有效索引。 2. 历史消息查询未分页或分页方式不对( LIMIT offset过深)。3. 离线消息拉取逻辑频繁扫表。 | 1. 使用EXPLAIN分析慢查询SQL。2. 监控数据库慢查询日志。 | 1. 确保conversation_id,receiver_id,status,created_at上有合适索引。2. 将 LIMIT offset分页改为基于id或created_at的游标分页。3. 离线消息尽量走Redis缓存,避免直接查库。 |
| 长连接数过多,网关内存溢出 | 1. 单机连接数达到上限。 2. 未及时清理失效连接。 | 1. 监控网关服务器的连接数、内存和CPU。 2. 检查是否有连接泄露(Channel未关闭)。 | 1.水平扩展网关:使用多个网关节点,客户端通过负载均衡连接。 2.强化连接管理:实现更精确的心跳和保活机制,定时扫描并踢掉无效连接。 3.优化Channel存储:使用更高效的数据结构,如 LongObjectHashMap。 |
| “已读”状态不同步 | 1. 客户端ACK丢失或未发送。 2. 多端登录时,一个端已读,状态未同步到其他端。 | 1. 检查客户端网络和ACK发送逻辑。 2. 查看服务端是否收到并处理了ACK。 | 1. 增强ACK可靠性,可加入重传机制。 2.已读状态同步:当某个端已读后,服务端应通过MQ广播一个 MESSAGE_READ事件,通知该用户的其他在线端更新本地状态。 |
7. 生产环境部署与监控建议
将系统投入生产环境,还需要考虑以下方面:
- 网关层集群化:使用 Nginx 的
ip_hash或基于userId的哈希负载均衡,将同一用户的路由到固定网关节点,便于连接查找。同时,网关节点需要无状态化,连接信息统一存于 Redis Cluster。 - 消息队列高可用:RabbitMQ 需配置镜像队列,Kafka 需配置多副本,确保消息不丢失。
- 数据库分库分表:当单表数据量过大时(如超过千万),需按
user_id或conversation_id进行分片。 - 全面的监控:
- 业务监控:消息发送量、送达率、已读率、在线用户数。
- 系统监控:各服务节点的CPU、内存、连接数、GC情况。
- 中间件监控:MySQL慢查询、Redis内存/命中率、MQ堆积情况。
- 链路追踪:集成 SkyWalking 或 Zipkin,追踪一条消息从发送到推送的完整路径,便于定位延迟瓶颈。
- 安全与限流:
- 鉴权:WebSocket 连接建立时必须进行强身份认证。
- 防刷:对发送消息的API进行频率限制。
- 内容安全:对消息内容进行敏感词过滤或图片鉴黄。
设计并实现一个站内信推送系统,是一个从“简单CRUD”思维迈向“分布式系统”思维的经典练习。它要求你综合考虑网络编程、数据存储、异步消息、缓存策略和实时通信。本文提供的架构与实现,是一个平衡了复杂度与功能的起点。你可以在此基础上,根据自身业务需求,引入更高级的特性,如消息撤回、消息编辑、端到端加密、富媒体消息、群聊等。记住,好的设计不是一步到位的,而是在清晰的架构约束下,随着业务演进不断迭代而成的。建议你先在本地环境跑通核心流程,再逐步将各个组件替换为集群化方案,最终构建出支撑亿级用户对话的可靠系统。