ARTICLE DETAIL

资讯详情

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

从零构建高可用站内信系统:WebSocket、消息队列与数据一致性实战

从零构建高可用站内信系统:WebSocket、消息队列与数据一致性实战 你是不是觉得站内信推送系统很简单不就是用户发个消息系统存一下然后推给另一个用户吗很多开发者一开始都这么想直到真正动手时才发现处处是坑消息延迟、已读状态不同步、海量数据下的性能瓶颈、推送失败如何补偿……一个看似简单的功能背后却涉及系统架构、数据一致性、实时性和扩展性等多个维度的设计挑战。这篇文章要解决的正是这个“看起来简单做起来头疼”的问题。我们将从一个真实的业务场景出发彻底拆解一个高可用、可扩展的站内信推送系统该如何设计与实现。读完本文你将不仅知道如何用代码实现收发消息更能掌握一套应对高并发、保证消息必达、支持多端同步的完整设计思路与工程实践。无论你是要为一个快速发展的社区、一个电商客服系统还是一个内部协作工具搭建消息通道这里都有你需要的答案。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) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT用户消息表;关键设计解析conversation_id这是优化查询的关键。将“查询用户A和用户B的所有消息”这个复杂条件(senderA AND receiverB) OR (senderB AND receiverA)简化为对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 groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-websocket/artifactId /dependency !-- 数据持久化 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-jpa/artifactId /dependency dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId scoperuntime/scope /dependency !-- 缓存与消息队列 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency !-- Netty (用于更底层的WebSocket实现可选) -- dependency groupIdio.netty/groupId artifactIdnetty-all/artifactId version4.1.108.Final/version /dependency !-- 工具类 -- dependency groupIdorg.projectlombok/groupId artifactIdlombok/artifactId optionaltrue/optional /dependency dependency groupIdcom.alibaba.fastjson2/groupId artifactIdfastjson2/artifactId version2.0.51/version /dependency /dependencies3.2 基础配置 (application.yml)# 文件路径src/main/resources/application.yml spring: datasource: url: jdbc:mysql://localhost:3306/message_db?useUnicodetruecharacterEncodingutf8useSSLfalseserverTimezoneAsia/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 ChannelInitializerSocketChannel() { 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 SimpleChannelInboundHandlerTextWebSocketFrame { 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 ApiResponseSendMessageResult 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 RedisTemplateString, 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; ListString 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 SetKey为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 JpaRepositoryUserMessage, 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) ListUserMessage 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”思维迈向“分布式系统”思维的经典练习。它要求你综合考虑网络编程、数据存储、异步消息、缓存策略和实时通信。本文提供的架构与实现是一个平衡了复杂度与功能的起点。你可以在此基础上根据自身业务需求引入更高级的特性如消息撤回、消息编辑、端到端加密、富媒体消息、群聊等。记住好的设计不是一步到位的而是在清晰的架构约束下随着业务演进不断迭代而成的。建议你先在本地环境跑通核心流程再逐步将各个组件替换为集群化方案最终构建出支撑亿级用户对话的可靠系统。
返回列表