ARTICLE DETAIL

资讯详情

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

SpringBoot+RabbitMQ+Redis+MySQL私信系统设计与高并发优化实践

SpringBoot+RabbitMQ+Redis+MySQL私信系统设计与高并发优化实践 私信模块在压测环境跑到500并发时MySQL连接池直接被打满报错日志里全是“Connection is not available, request timed out”。这是我接手社交平台私信系统重构时遇到的第一道坎。如果你也在做私信、站内信、一对一消息这类功能大概率会踩到同样的坑消息量级看着不大读写路径却远比想象中复杂。这篇文章我会把整套基于 SpringBoot RabbitMQ Redis MySQL 的私信系统设计思路、核心代码、表结构、以及上线前后踩过的坑完整拆开讲适合正在规划私信模块、或者想优化现有聊天类功能的同学参考。1. 私信系统为什么不是“聊天框”那么简单需求拆解与四大件分工1.1 私信与IM、聊天室的本质差异很多做业务系统的人一开始会把私信想简单了觉得无非是“A发一条消息给BB能看到”然后顺手就把消息直接写进一张表再用轮询或者WebSocket推一下。真正做起来才发现私信系统和即时通讯软件IM、聊天室有着完全不同的技术侧重。IM 的核心诉求是极低延迟和长连接保活所以业界常见方案是自研协议栈、Netty长连接、私有推送通道。聊天室的核心诉求是消息广播和在线人数承载大量使用发布订阅模型。而私信系统的核心诉求是什么是不丢消息、可追溯、会话隔离、多端同步。也就是说私信用户不会要求你在50毫秒内把消息送到对方手里晚个一两秒问题不大但如果你把消息弄丢了或者对方换个设备登录后发现历史消息对不上那就是事故。所以私信系统在设计时可靠性、可回溯性、数据一致性要优先于实时性。另一个现实约束是大多数团队的私信系统是嵌在已有社交平台里的用户体系、好友关系、内容审核这些模块都已经存在。你不能为了一个私信功能就自建一套IM基础设施必须复用现有的业务架构。这就决定了技术选型要往“工程化、可维护、能复用已有组件”的方向走而不是从零堆轮子。1.2 选型逻辑为什么是 SpringBoot RabbitMQ Redis MySQL先说结论这个组合不是性能最优解但它是业务系统里性价比最高、团队最容易维护的组合。组件核心职责如果不用它会发生什么SpringBoot接口层、业务编排、参数校验、鉴权开发效率大幅下降与其他业务模块整合成本高RabbitMQ异步削峰、消息投递、解耦生产者和消费者发送链路全部同步数据库在高并发下被击穿Redis在线状态、未读数、会话列表热数据、幂等去重MySQL读压力过大热点会话查询慢MySQL会话数据、历史消息、多端同步的唯一数据源消息无法可靠持久化后续功能扩展受限这里我想重点说说为什么不用Netty。很多做IM的人会推荐直接上Netty但它意味着你要自己处理编解码、心跳、断线重连、消息序号、ACK确认、离线消息补偿等一整套IM协议层问题。对于业务系统里的私信模块来说这些工作量和风险是不划算的。实际上私信的实时推送完全可以用 WebSocket SpringBoot 来实现再配合RabbitMQ做消息分发生产环境足够稳。MySQL在这个架构里不是被替代而是作为冷热数据的最终归属地。私信消息一旦确认送达就必须落到MySQL里面它是用户多端同步和消息追溯的基础。1.3 消息从发送到展示的完整链路我用文字把整条链路模型交代一下不画图你跟着顺序读就能在脑子里构建出结构用户A调用发送接口请求进入 SpringBoot 的 Controller 层Service 层做参数校验、会话关系校验、敏感词过滤、频控检查通过后把消息内容、发送者、接收者、会话ID组装成消息对象消息先写Redis会话列表、未读数变更再投递到RabbitMQRabbitMQ 的消费者异步消费消息写入 MySQL如果接收方在线通过 WebSocket 实时推送如果离线等他下次上线时拉取接收方已读消息后反向更新Redis和MySQL中的已读状态注意到第4步的顺序没有先缓存、后队列、再落库。这是这套设计里最核心的一个决策后面会展开讲。2. 发送链路改造从同步写库到RabbitMQ异步削峰2.1 压测暴露的写路径瓶颈我接手的时候老系统的发送逻辑非常简单粗暴发送接口收到请求后直接往 private_message 表插一条记录然后 update 会话表的 last_message 字段最后返回成功。从功能上看没什么毛病但压测数据出来的时候大家都沉默了——500并发下平均响应时间超过3秒MySQL连接池被打满一堆请求直接超时。问题出在哪同步写库的性能瓶颈非常直接一次请求要完成一次 INSERT 和一次 UPDATE两个操作都落在同一台MySQL上。500并发意味着每秒有上千次写操作InnoDB的刷盘能力和行锁竞争根本扛不住。更扎心的是用户发一条消息他并不需要立刻看到“已落库”的永久性结果他只需要看到“消息已发送”状态置为发送中就够了。所以发送链路的第一刀就是把写MySQL的操作从同步调用里挪出去。2.2 发送接口只做三件事校验、缓存、投递重构之后发送接口的职责边界被压缩得很清爽校验、更新缓存、投递MQ。核心Service方法大概是这样的Transactional public SendResult sendPrivateMessage(SendMessageRequest request) { // 1. 校验用户是否存在、双方是否允许私信、消息内容是否合规 checkSendPermission(request.getFromUid(), request.getToUid()); checkContentRisk(request.getContent()); // 2. 组装消息对象生成全局唯一消息ID雪花算法 PrivateMessage message buildMessage(request); // 3. 先更新Redis会话列表排序、未读数自增 refreshConversationCache(message); incrUnreadCount(message.getToUid(), message.getConversationId()); // 4. 再投递到RabbitMQ落库由消费者异步完成 sendToMq(message); // 5. 返回发送成功此时消息状态是“发送中” return SendResult.success(message); }注意这里的Transactional实际上管理的是Redis操作的事务边界而不是数据库事务。这就是关键的一个转变发送结果不再依赖数据库的实时写入而是依赖缓存更新和MQ投递成功。同步调用时间从原来的几十毫秒、遇到高峰期几百毫秒甚至超时压到了稳定的10毫秒以下。用户体感就是“发消息点了就发出去了”没有再转圈等待。2.3 RabbitMQ配置交换机、队列、路由键的设计RabbitMQ部分很容易被忽略的就是vhost和权限隔离。很多生产环境里多个团队共用同一个RabbitMQ中间件如果不做vhost隔离消息队列之间互相干扰、配置冲突是家常便饭。我的建议是私信系统单独创建一套虚拟主机并单独建账号只赋予这个vhost的操作权限。队列名称也带上业务前缀。spring: rabbitmq: host: 127.0.0.1 port: 5672 username: private_chat password: xxxx virtual-host: /private_chat publisher-confirm-type: correlated publisher-returns: true交换机类型选择上我一般用topic 交换机这样路由规则最灵活。Bean public TopicExchange privateChatExchange() { return new TopicExchange(private.chat.exchange, true, false); } Bean public Queue privateChatQueue() { return QueueBuilder.durable(private.chat.queue).build(); } Bean public Binding binding() { return BindingBuilder.bind(privateChatQueue()) .to(privateChatExchange()) .with(private.chat.#); }生产者投递时必须开启确认模式确保消息真正到达了BrokerrabbitTemplate.convertAndSend( private.chat.exchange, private.chat. conversationId, message, correlationData // CorrelationData 里带上消息ID用于回调确认 );2.4 消费者手动ACK、幂等与死信兜底消费者这段是整个发送链路里坑最多的部分。默认的自动ACK模式在消费端处理失败时会丢消息所以必须改成手动ACK。消费逻辑遵循一个固定套路先做幂等校验再写库最后ACK如果前置校验失败直接ACK但不处理如果写库失败不ACK让消息重回队列或者进入死信队列。RabbitListener(queues private.chat.queue) public void onMessage(Message message, Channel channel) throws IOException { long deliveryTag message.getMessageProperties().getDeliveryTag(); PrivateMessage msg JSON.parseObject(message.getBody(), PrivateMessage.class); // 幂等校验消息ID是否已消费过 Boolean firstConsume redisTemplate.opsForValue() .setIfAbsent(msg:consume: msg.getMessageId(), 1, Duration.ofDays(7)); if (!Boolean.TRUE.equals(firstConsume)) { channel.basicAck(deliveryTag, false); return; } try { privateMessageMapper.insert(msg); conversationMapper.updateLastMessage(msg); channel.basicAck(deliveryTag, false); } catch (Exception e) { channel.basicNack(deliveryTag, false, true); // requeue true } }这里有一个细节basicNack的 requeue 参数如果设为 true消息会回到队列头部附近导致消费失败的消息反复被同一个消费者拉取形成消息堆积的心跳陷阱。更好的做法是 requeuefalse把消费失败的消息转入死信队列再由单独的死信消费者做补偿处理。另外一个容易踩的点是RabbitMQ 消息重试没有内置的退避机制如果消费者代码有问题重试会非常高频。所以在消费端代码里要加一个简单的延迟重试控制或者直接用Spring Retry做指数退避避免消息故障时把消费者线程池拖垮。3. Redis在私信系统里的三块硬骨头会话列表、未读数与在线状态3.1 会话列表的缓存建模用户打开私信页面时第一个看到的是会话列表按最后一条消息的时间倒序排列。这个接口的数据库查询其实挺重需要 join 会话表、消息表和用户信息表。在高频访问下性能完全依赖缓存。我用Redis的 ZSET 来存储每个用户的会话列表key 是conversation:list:{userId}member 是 conversationIdscore 是最后一条消息的时间戳。ZSET天然支持按时间排序新增一条消息时只需要zadd一下就能保证列表排序正确。public void refreshConversationCache(PrivateMessage message) { String keyA conversation:list: message.getFromUid(); String keyB conversation:list: message.getToUid(); double score (double) message.getCreateTime().getTime(); redisTemplate.opsForZSet().add(keyA, message.getConversationId(), score); redisTemplate.opsForZSet().add(keyB, message.getConversationId(), score); // 设置过期时间避免长期不活跃用户的key占用内存 redisTemplate.expire(keyA, Duration.ofDays(30)); redisTemplate.expire(keyB, Duration.ofDays(30)); }会话列表接口只需要从ZSET里取前20个会话ID再批量从缓存或MySQL回源会话元数据。这里要注意ZSET里存了会话ID但会话对端的头像、昵称这些信息经常变化不应该冗余到Redis里应该由接口层去用户服务批量查询。把会话ID和用户信息分开存缓存一致性会好维护很多。3.2 未读数的读写路径未读数这个字段看起来简单做起来最大的问题是“一定要准”。未读数的处理我采用了先写Redis、异步落库MySQL的方案。发送方发消息时对接收方执行incrredisTemplate.opsForValue().increment(unread: toUid : conversationId);接收方打开会话时把该会话的未读数清零并异步通知后端public void readConversation(Long userId, Long conversationId) { String key unread: userId : conversationId; Integer unread (Integer) redisTemplate.opsForValue().get(key); if (unread ! null unread 0) { redisTemplate.delete(key); // 异步通知消费者把已读状态同步到MySQL sendReadAckToMq(userId, conversationId); } }为什么整体未读数不直接存一个总数因为用户在不同会话上需要分别显示红色角标所以必须按会话维度存储。另外总未读数可以设置一个key再单独incr但一次incr操作和一个hincrby操作的性能差别不大主要的坑在于清理时机一定要在“全部会话已读”时才允许清零总未读数否则角标会不对。3.3 在线状态与多端推送在线状态我用的是 Redis String 结构key 为online:user:{userId}value 里记录连接节点ID和最后心跳时间过期时间设为90秒。用户通过WebSocket连上来时写一次之后每30秒发一次心跳续期。有了在线状态之后消息推送策略就是典型的分支判断接收方在线立即通过WebSocket推送接收方离线则等待上线拉取。WebSocket推送到代码怎么拿到连接因为受限于每台服务器只能处理落在自己节点上的连接必须做一个轻量的路由连接建立时把 userId 和节点信息写入Redis消息要推送时查询用户连接的节点通过节点内部的本地连接池发送或者走服务间调用转发。如果只有一个节点这部分可以简化。Redis的可视化排查工具推荐用 Another Redis Desktop Manager这种客户端能直接看key的TTL和内存占用排查未读数key堆积的问题非常方便。3.4 缓存穿透、热点Key与大V私信私信场景里最容易出现的缓存问题有两个一个是缓存穿透另一个是热点Key。缓存穿透发生在会话列表突然被大量请求时如果Redis里查不到就回源MySQL但MySQL里也没有数据导致每次请求都打到数据库。解决方式一般为布隆过滤器或者空值缓存。私信场景中会话ID通常有规律我建议直接对不存在的会话ID写一个短暂的空值缓存TTL 30秒逻辑简单有效。热点Key则更头疼。假设平台里有个几百万粉丝的大V他的私信列表在短时间内涌入大量消息所有粉丝发消息时都要操作conversation:list:{大V用户ID}和unread:{大V用户ID}:{会话ID}这些key就成了热点。我看到过整个Redis节点CPU被打满的情况。应对方式有两个方向一是给热点用户的未读计数拆多个分片key比如unread:{userId}:{partition}读取时求和二是消息投递时不直接incr而是让MQ消费者在异步阶段批量合并更新。后一种方式对吞吐提升更明显因为消费者可以按用户分组攒够100条或者100毫秒窗口后再批量更新一次Redis减少热点写操作数量。4. MySQL存储层的表设计与查询治理4.1 核心表结构会话表 消息表Redis承担了热数据的读写MySQL则负责冷数据的可靠存储。表结构设计直接影响后续的消息追溯、多端同步和扩展。实际项目中我会设计两张核心表第一张是私信会话表CREATE TABLE private_conversation ( id bigint NOT NULL AUTO_INCREMENT COMMENT 会话ID, user_id bigint NOT NULL COMMENT 会话归属用户ID, peer_id bigint NOT NULL COMMENT 会话对方用户ID, last_message_id bigint DEFAULT NULL COMMENT 最后一条消息ID, last_message_at datetime DEFAULT NULL COMMENT 最后消息时间, unread_count int NOT NULL DEFAULT 0 COMMENT 未读数, created_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (id), KEY idx_user_updated (user_id, last_message_at), UNIQUE KEY uk_user_peer (user_id, peer_id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT私信会话表;这里最重要的设计是一个会话关系存储了两行即 (user_id A, peer_id B) 和 (user_id B, peer_id A)。这样每个用户查询自己的会话列表只需要按user_id索引扫描不需要额外的关联关系判断。第二张是消息明细表CREATE TABLE private_message ( id bigint NOT NULL AUTO_INCREMENT COMMENT 消息ID, conversation_id bigint NOT NULL COMMENT 会话ID, sender_id bigint NOT NULL COMMENT 发送者ID, receiver_id bigint NOT NULL COMMENT 接收者ID, content_type tinyint NOT NULL DEFAULT 1 COMMENT 内容类型 1文本 2图片 3语音, content text COMMENT 消息内容, status tinyint NOT NULL DEFAULT 1 COMMENT 状态 1正常 2撤回, created_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (id), KEY idx_conversation_time (conversation_id, id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT私信消息明细表;消息表主键用自增id配合(conversation_id, id)联合索引能覆盖“查询某个会话的历史消息”这个核心路径。为什么不直接用雪花ID做主键因为InnoDB的聚簇索引对自增值顺序插入更友好而雪花ID作为二级索引即可。4.2 游标分页与PageHelper的局限性很多团队做分页时习惯直接用 PageHelper 这类MyBatis分页插件在私信历史消息查询上我强烈建议不要这么做。原因很简单PageHelper 基于 LIMIT OFFSET 实现随着翻页深度增加offset 越来越大MySQL需要扫描前面所有不需要的行再丢弃性能急剧下降。在消息表这种数据量持续增长的场景里翻到第100页时SQL可能已经慢到几百毫秒了。正确做法是使用游标分页即基于上一页最后一条消息ID做条件查询public ListPrivateMessage pageHistory(Long conversationId, Long lastMsgId, Integer limit) { QueryWrapperPrivateMessage wrapper new QueryWrapper(); wrapper.eq(conversation_id, conversationId) .and(w - w.lt(id, lastMsgId)) // 上一页最后一条消息ID .orderByDesc(id) .last(LIMIT limit); return privateMessageMapper.selectList(wrapper); }客户端传入 lastMsgId 即可第一次查时传一个很大的值。这种分页方式无论翻多深走的都是(conversation_id, id)索引性能稳定。4.3 会话列表回源与更新时机会话列表从ZSET读取但Redis里的数据只保存最后消息时间、会话ID这些轻量信息会话的展示名、对方头像等信息需要回源MySQL或者用户服务。这里有一个明显的钩子如果消息发出去后异步落库失败Redis里能看到会话但MySQL里没有对应的新消息那多端同步就会出现数据空洞。所以异步落库是核心链路的必经之路为了不让MySQL成为瓶颈写库操作不搞实时join查询而是直接插入消息记录然后 update 会话表的last_message_at。每次更新会在idx_user_updated索引上顺序追加性能可控。4.4 数据归档与冷热分离私信消息是持续增长的全部留在主表里总有一天会把MySQL拖垮。我建议按时间做冷热分离90天内的消息放主表超过90天的定期迁移到归档表。因为私信的特点是绝大部分用户只会看近期消息90天前的历史消息被主动翻出来的概率很低。归档任务可以每天凌晨执行一次用批量查询加批量插入的方式迁移迁移完成后按时间范围删除原纪录。注意一次别删太多否则会产生大批量行锁和binlog文件膨胀我的经验是一批1000条执行完sleep一下。同时消息表的最佳实践是给created_at建索引否则归档任务的WHERE条件会变成全表扫描。5. 生产环境必踩的坑从RabbitMQ权限到SpringBoot版本兼容5.1 Docker部署RabbitMQ后admin账号无法创建虚拟主机这个坑几乎每个人都会踩一遍。用Docker起了一个RabbitMQ管理界面能打开但是用admin账号登录后发现无论如何都没法创建虚拟主机页面一直提示权限不足。根因是RabbitMQ默认的guest用户只有localhost访问权限而admin账号即使创建出来了如果没有被赋予对应vhost的操作权限它什么都做不了。这就好比你有登录密码但没有钥匙开里面的门。解决办法是在容器里执行命令把vhost权限完整授权给admindocker exec -it rabbitmq rabbitmqctl add_vhost /private_chat docker exec -it rabbitmq rabbitmqctl set_permissions -p /private_chat admin .* .* .*先用rabbitmqctl list_vhosts确认当前有哪些虚拟主机再用rabbitmqctl list_permissions -p /private_chat检查授权。凡是遇到账号能登录但不能操作的情况先查权限列表十有八九是权限没配。5.2 SpringBoot版本“太高”引发的适配问题做这套私信系统时我正赶上Spring Boot 3.x普及很多依赖的兼容问题都被网友吐槽过。最大的变化是命名空间从javax.*换成了jakarta.*导致老项目的代码导入直接编译报错。还有一些底层实现调整比如Spring Boot 3.x默认使用Caffeine做本地缓存如果老代码里用Guava Cache就要留意缓存行为的变化。另外Spring Boot版本升级后RabbitMQ的连接池配置项也变了。以前用的spring.rabbitmq.cache.channel.size在3.x里作用不再明显新版连接复用的机制不同需要显式调整连接工厂的并发参数。建议升级Spring Boot大版本时把中间件相关配置按官方文档重新扒一遍不要只图编译通过。5.3 消息乱序问题与多消费者并发私信这种一对一会话里消息乱序是用户能明显感知的Bug。场景是这样的消息消费者部署了多个实例RabbitMQ会把同一个会话的相邻两条消息分发给不同的消费者实例执行由于线程调度差异后发的那条消息可能先完成落库用户刷新页面时就看到两条消息的顺序反了。说话顺序是私信的底线不能乱。解决办法是把同一会话的消息路由到同一个消费者实例我用的是一致性哈希交换机思路让相同conversationId的消息固定走到同一个队列。RabbitMQ从3.x之后支持一致性哈希交换机但版本较老则需要手动实现。更简单的做法是队列分区预先创建N个队列比如16个根据conversationId % N决定投递到哪个队列每个消费者实例固定只绑定一个队列。这样同一会话的消息永远不会被并发消费顺序天然有保证。缺点是一个实例挂了以后它消费的那几个队列的消息会停止处理需要靠死信和补偿机制兜底在普通业务量下利大于弊。5.4 消费积压的监控与排查清单私信系统上线后RabbitMQ消费积压是必须实时盯的指标。我有几个长期有效的排查手段RabbitMQ管理后台重点看ready和unacknowledged两个数字ready长期大于0说明有积压。Redis慢日志SLOWLOG GET要定期看未读数热点key导致CPU飙升时慢日志能帮快速定位。MySQL慢查询日志打开对执行时间超过100ms的SQL做分析。私信场景最常出现的慢SQL十有八九是会话列表查询里做了大偏移量分页或者缺失索引的join。有一段时间我发现消费者消费速率骤降排查了半天最后发现是消息队列里的消息体太大——用户发了一段超长文本超过1MBRabbitMQ默认的消息大小限制直接触发了拒绝。这个问题遇到一次就记住了定制化业务场景下一定要在发送接口里限制单条消息长度比如文本不超过5000字图片用url代替base64。最终的一个小建议与个人体会这套私信系统重构上线后最直观的变化是发送接口的P99耗时从800ms降到了120ms左右MySQL连接数从打满状态稳定在几十。我个人的体会是私信这类业务真正的技术难点从来不是“消息怎么发出去”而是数据在多个中间件之间流转时如何保证不丢、不乱、不重复。如果你要在自己的系统里落地这套架构建议先画清楚一条消息从发送到展示的完整时序图和异常分支然后再动手写代码。上线前一定压测压测时不仅看吞吐还要盯着RabbitMQ的积压量、Redis的OPS和MySQL的连接数三个指标。最后再分享一个我在实际排查过程中觉得特别管用的小工具链RabbitMQ的rabbitmqctl list_queues name messages messages_ready messages_unacknowledged一条命令能快速定位积压情况Redis用redis-cli --bigkeys扫描大key能避免热点key把节点内存打满。这些小操作在关键时刻比任何一个复杂的监控系统都管用。
返回列表